feat: add Grafana-ready SQLite observability (#86)
* feat: add Grafana-ready SQLite observability (#79, #80, #81) New `mt5cli/grafana.py` module with idempotent DDL helpers: - `create_snapshot_tables` — five SQLite tables for time-series account, position, order, terminal, and run-status snapshots - `create_grafana_views` — 13 `grafana_*` views with integer epoch-second `time` columns; missing source tables emit warnings and are skipped - `create_grafana_indexes` — 9 performance indexes guarded by column checks - `ensure_grafana_schema` — convenience wrapper calling all three above - Insert helpers: `insert_account_snapshot`, `insert_position_snapshots`, `insert_order_snapshots`, `insert_terminal_snapshot`, `record_snapshot_run` New stable SDK exports in `mt5cli.__init__` and `mt5cli.contract`: - `update_observability` — appends a timestamped snapshot to a SQLite db using an already-connected `Mt5DataClient`; never places orders - `update_observability_with_config` — standalone wrapper that opens and closes the MT5 connection automatically New CLI commands (Collection panel): - `grafana-schema` — idempotent schema setup, no MT5 connection required - `snapshot` — append account/position/order/terminal rows; supports `--symbol`, `--with-account/--no-account`, and equivalent flags All public modules maintain 100 % branch coverage; 968 tests pass. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * chore: apply mdformat to docs after grafana observability additions Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * chore: bump version to 1.1.0 Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: address PR #86 review feedback - Fix unaggregated time in grafana_realized_pnl GROUP BY query (MAX) - Replace O(N) per-symbol API calls with single call + client-side filter - Eliminate double create_snapshot_tables when with_grafana_schema=True - Move grafana imports to module level in sdk.py (remove PLC0415 noqa) - Default with_grafana_schema to False (run grafana-schema once for setup) - Fix README position example to use snapshot_runs for latest snapshot - Update tests to reflect new behavior and correct patch targets Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * chore: apply ruff formatting and sync lock file Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: address owner review feedback on PR #86 - Filter grafana_*_snapshots views to only expose rows from successful runs (JOIN snapshot_runs WHERE status='ok'), closing the partial-snapshot visibility gap raised in PRRT_kwDORzI_286MvDJ6 - Add issubset column guards for snapshot table index creation, consistent with the rest of create_grafana_indexes (PRRT_kwDORzI_286MvUx1) - Fix README example queries: views expose 'time' not 'observed_at' (PRRT_kwDORzI_286MvUxw) - Correct public-contract.md default for with_grafana_schema (False, not True) and point to grafana-schema command (PRRT_kwDORzI_286MvUxy) Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: add run_id to snapshot schema and fix stale-positions README query Replace second-level observed_at as the join key between snapshot_runs and snapshot tables with a stable run_id INTEGER PRIMARY KEY. Two runs in the same second now get distinct run_ids, preventing view duplication and cross-contamination from a failed run. Update README example to use snapshot_runs for latest-snapshot lookup so zero-position runs return an empty result instead of stale rows. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * chore: format markdown tables Align table column widths in README and public-contract documentation. Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com> * fix: close sqlite connections deterministically * fix: require entry filter in grafana_realized_pnl and add time to grafana_trade_stats grafana_realized_pnl now requires the entry column and filters to close-side deals (entry IN (1, 2, 3)), consistent with grafana_symbol_pnl and grafana_trade_stats. grafana_trade_stats now requires the time column and emits MAX(time_expr) AS "time" so it satisfies the documented Grafana view contract (integer epoch-second time column throughout). Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: normalize pd.Timestamp time_setup to epoch int in insert_order_snapshots orders_get_as_df() returns datetime-converted columns by default, so time_setup is a pd.Timestamp in normal use. Passing it directly to sqlite3.executemany raises ProgrammingError. Added _to_epoch_int helper that converts datetime.datetime subclasses (including pd.Timestamp) and raw int/float values to integer epoch seconds, returning None for other types. Regression tests cover all four input paths. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: pre-drop all grafana_* views at start of create_grafana_views Previously, a builder that skipped due to a missing source table or column would not drop the view it owned, leaving stale views referencing gone tables. Now create_grafana_views drops all 13 known grafana_* views before calling any builder, so a schema refresh always removes views whose source has disappeared. Regression test covers the create → drop-source → refresh cycle. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: expose run_id in snapshot views and drop time from summary views Grafana snapshot views now expose run_id so latest-state queries can use MAX(run_id) instead of the ambiguous second-level MAX(observed_at). grafana_realized_pnl and grafana_trade_stats lose their MAX(time) column and are reclassified as static summary views; their all-time aggregates are not filterable by Grafana time-range selectors. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> * fix: guard snapshot views against missing run_id and fix README view docs _build_snapshot_view now skips with a warning when the underlying snapshot table exists but lacks a run_id column, preventing a broken view that fails at query time. Adds a regression test for that path. README Grafana section now qualifies that grafana_realized_pnl and grafana_trade_stats are static summary views (no time column) and splits the view table to match docs/api/public-contract.md. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com> --------- Co-authored-by: agent <agent@localhost> Co-authored-by: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -185,6 +185,8 @@ python -m mt5cli -o account.csv account-info
|
||||
| `order-send` | Send a raw trade request to the trade server (`--yes` required; expert path) |
|
||||
| `close-positions` | Close open positions by `--symbol` or `--ticket` (`--yes` required for live; `--dry-run` available) |
|
||||
| `collect-history` | Collect rates, history-orders, and history-deals for one or more symbols into a single SQLite database (ticks opt-in via `--dataset ticks`) |
|
||||
| `grafana-schema` | Create or refresh Grafana-ready views and indexes in an existing SQLite database (idempotent, no MT5 connection) |
|
||||
| `snapshot` | Snapshot current account, position, order, and terminal state into SQLite for live Grafana dashboards |
|
||||
|
||||
Use `order-check` to validate a request payload before running `order-send --yes`.
|
||||
`close-positions` is the safer high-level alternative that builds correct close
|
||||
@@ -204,6 +206,104 @@ 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 `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.
|
||||
|
||||
### Grafana-ready SQLite dashboards
|
||||
|
||||
mt5cli can prepare a SQLite database for use as a Grafana datasource (via the [SQLite plugin](https://grafana.com/grafana/plugins/frser-sqlite-datasource/) or similar). Most `grafana_*` views expose an integer epoch-second `time` column for use in Grafana time-series panels. Two views (`grafana_realized_pnl`, `grafana_trade_stats`) are static symbol-level summaries with no `time` column — use them in table or stat panels.
|
||||
|
||||
#### Prepare the schema (idempotent, no MT5 connection needed)
|
||||
|
||||
```bash
|
||||
mt5cli -o history.db grafana-schema
|
||||
```
|
||||
|
||||
This creates snapshot tables (`account_snapshots`, `position_snapshots`, `order_snapshots`, `terminal_snapshots`, `snapshot_runs`) and all `grafana_*` views and indexes in the SQLite database. Safe to run repeatedly — all operations are idempotent.
|
||||
|
||||
#### Snapshot current account state
|
||||
|
||||
```bash
|
||||
mt5cli -o history.db snapshot \
|
||||
--symbol JP225 --symbol HK50 --symbol NL25 \
|
||||
--with-account --with-positions --with-orders --with-terminal \
|
||||
--with-grafana-schema
|
||||
```
|
||||
|
||||
Appends one timestamped row per data type. Never places orders or modifies trading state. Run periodically (e.g. from a cron job or a loop) to build a time-series account history.
|
||||
|
||||
#### SDK usage
|
||||
|
||||
```python
|
||||
from pdmt5 import Mt5DataClient, Mt5Config
|
||||
from mt5cli import update_observability, update_observability_with_config
|
||||
|
||||
# Reuse an already-connected client
|
||||
client = Mt5DataClient(config=Mt5Config(login=12345))
|
||||
client.initialize_and_login_mt5()
|
||||
try:
|
||||
update_observability(
|
||||
client=client,
|
||||
output="history.db",
|
||||
symbols=["EURUSD", "GBPUSD"], # optional position/order filter
|
||||
include_account=True,
|
||||
include_positions=True,
|
||||
include_orders=True,
|
||||
include_terminal=True,
|
||||
with_grafana_schema=True,
|
||||
)
|
||||
finally:
|
||||
client.shutdown()
|
||||
|
||||
# Standalone wrapper that opens/closes MT5 automatically
|
||||
update_observability_with_config(
|
||||
output="history.db",
|
||||
config=Mt5Config(login=12345),
|
||||
)
|
||||
```
|
||||
|
||||
#### Available Grafana views
|
||||
|
||||
**Time-series views** (integer epoch-second `time` column; snapshot views also expose `run_id`):
|
||||
|
||||
| View | Source | Description |
|
||||
| ---------------------------- | -------------------- | ---------------------------------------------------------- |
|
||||
| `grafana_rates` | `rates` | OHLCV bars with integer epoch `time` |
|
||||
| `grafana_ticks` | `ticks` | Tick data with integer epoch `time` |
|
||||
| `grafana_history_deals` | `history_deals` | All deals with epoch `time` |
|
||||
| `grafana_history_orders` | `history_orders` | All historical orders; adds epoch `time` from `time_setup` |
|
||||
| `grafana_trade_deals` | `history_deals` | Trade deals only (`type IN (0,1)`) |
|
||||
| `grafana_cash_events` | `history_deals` | Non-trade deals (deposits, dividends, etc.) |
|
||||
| `grafana_symbol_pnl` | `history_deals` | Per-close-deal profit/loss per symbol |
|
||||
| `grafana_account_snapshots` | `account_snapshots` | Account balance/equity/margin time series |
|
||||
| `grafana_position_snapshots` | `position_snapshots` | Open position snapshots over time |
|
||||
| `grafana_order_snapshots` | `order_snapshots` | Active order snapshots over time |
|
||||
| `grafana_terminal_snapshots` | `terminal_snapshots` | Terminal connectivity snapshots |
|
||||
|
||||
**Static summary views** (no `time` column; use in table or stat panels, not time-series):
|
||||
|
||||
| View | Source | Description |
|
||||
| ---------------------- | --------------- | ------------------------------------- |
|
||||
| `grafana_realized_pnl` | `history_deals` | Cumulative realized PnL per symbol |
|
||||
| `grafana_trade_stats` | `history_deals` | Win/loss counts and profit per symbol |
|
||||
|
||||
#### Example Grafana queries
|
||||
|
||||
```sql
|
||||
-- Equity curve over time
|
||||
SELECT time, equity FROM grafana_account_snapshots ORDER BY time;
|
||||
|
||||
-- Rolling balance by account login
|
||||
SELECT time, login, balance FROM grafana_account_snapshots
|
||||
WHERE login = $login ORDER BY time;
|
||||
|
||||
-- Open positions at latest successful snapshot
|
||||
SELECT symbol, volume, profit FROM grafana_position_snapshots
|
||||
WHERE run_id = (SELECT MAX(run_id) FROM snapshot_runs WHERE status = 'ok');
|
||||
|
||||
-- Realized PnL by symbol
|
||||
SELECT symbol, total_profit FROM grafana_trade_stats ORDER BY total_profit DESC;
|
||||
```
|
||||
|
||||
> **Note**: OpenTelemetry integration is intentionally not part of this release and is tracked separately.
|
||||
|
||||
### Incremental history SDK
|
||||
|
||||
For automated pipelines, use the importable incremental API instead of re-fetching fixed date ranges:
|
||||
|
||||
@@ -124,6 +124,63 @@ sending requests. Failed, malformed, or unknown broker retcodes are fail-closed
|
||||
and returned as `status="failed"` with normalized `request` / `response` details;
|
||||
`dry_run=True` never calls `ensure_symbol_selected()` or `order_send()`.
|
||||
|
||||
### Grafana observability (SQLite read model)
|
||||
|
||||
These helpers prepare a SQLite database as a Grafana datasource. All DDL is
|
||||
idempotent (`CREATE TABLE IF NOT EXISTS`, `DROP VIEW IF EXISTS` + `CREATE
|
||||
VIEW`, `CREATE INDEX IF NOT EXISTS`). Missing source tables are skipped with a
|
||||
warning rather than raising an error.
|
||||
|
||||
| Symbol | Role |
|
||||
| ---------------------------------- | ----------------------------------------------------------------------------------------------- |
|
||||
| `update_observability` | Append one timestamped snapshot row per data type; accepts an already-connected `Mt5DataClient` |
|
||||
| `update_observability_with_config` | Standalone wrapper: opens/closes MT5 connection automatically around `update_observability` |
|
||||
|
||||
Both functions write to the SQLite path given by `output=`. The optional
|
||||
`symbols` parameter filters `positions_get` / `orders_get` by symbol.
|
||||
`with_grafana_schema=False` (default) skips Grafana view/index setup; run
|
||||
`grafana-schema` once to set up the schema, then call `snapshot` repeatedly
|
||||
without this flag.
|
||||
|
||||
**Snapshot tables** (created by `create_snapshot_tables` in `mt5cli.grafana`):
|
||||
|
||||
| Table | Content |
|
||||
| -------------------- | ----------------------------------------- |
|
||||
| `account_snapshots` | Balance, equity, margin, free-margin, P&L |
|
||||
| `position_snapshots` | Open positions: symbol, volume, profit, … |
|
||||
| `order_snapshots` | Active orders: symbol, type, price, … |
|
||||
| `terminal_snapshots` | Terminal connectivity and build info |
|
||||
| `snapshot_runs` | Per-run status (`ok` / `error`) timestamp |
|
||||
|
||||
**Grafana time-series views** (integer epoch-second `time` column; snapshot views also expose `run_id`):
|
||||
|
||||
| View | Source |
|
||||
| ---------------------------- | -------------------------------- |
|
||||
| `grafana_rates` | `rates` table |
|
||||
| `grafana_ticks` | `ticks` table |
|
||||
| `grafana_history_deals` | `history_deals` |
|
||||
| `grafana_history_orders` | `history_orders` |
|
||||
| `grafana_trade_deals` | `history_deals` trade types only |
|
||||
| `grafana_cash_events` | `history_deals` non-trade events |
|
||||
| `grafana_symbol_pnl` | Per-close-deal P&L per symbol |
|
||||
| `grafana_account_snapshots` | `account_snapshots` |
|
||||
| `grafana_position_snapshots` | `position_snapshots` |
|
||||
| `grafana_order_snapshots` | `order_snapshots` |
|
||||
| `grafana_terminal_snapshots` | `terminal_snapshots` |
|
||||
|
||||
**Grafana static summary views** (no `time` column; use for table/stat panels, not time-series):
|
||||
|
||||
| View | Source |
|
||||
| ---------------------- | ------------------------------------- |
|
||||
| `grafana_realized_pnl` | Cumulative realized PnL per symbol |
|
||||
| `grafana_trade_stats` | Win/loss counts and profit per symbol |
|
||||
|
||||
Lower-level helpers (`ensure_grafana_schema`, `create_grafana_views`,
|
||||
`create_grafana_indexes`, `create_snapshot_tables`, `start_snapshot_run`,
|
||||
`insert_account_snapshot`, `insert_position_snapshots`, `insert_order_snapshots`,
|
||||
`insert_terminal_snapshot`, `record_snapshot_run`) are available directly from
|
||||
`mt5cli.grafana` and are not part of the package-root stable surface.
|
||||
|
||||
### Errors
|
||||
|
||||
| Symbol | Role |
|
||||
@@ -135,14 +192,15 @@ and returned as `status="failed"` with normalized `request` / `response` details
|
||||
Lower-level helpers are available from their owning modules and are not part
|
||||
of the package-root stable surface. Import them directly when needed:
|
||||
|
||||
| Module | Examples |
|
||||
| ------------------- | ---------------------------------------------------------------------------------------------- |
|
||||
| `mt5cli.history` | `resolve_rate_view_name`, `resolve_rate_tables`, `load_rate_data`, `build_rate_view_name` |
|
||||
| `mt5cli.sdk` | `copy_rates_from`, `copy_ticks_from`, `account_info`, `symbols`, `mt5_summary`, `latest_rates` |
|
||||
| `mt5cli.schemas` | `DataKind`, `normalize_dataframe`, `validate_schema`, `DEDUP_KEYS` |
|
||||
| `mt5cli.utils` | `Dataset`, `IfExists`, `detect_format`, `export_dataframe`, `export_dataframe_to_sqlite` |
|
||||
| `mt5cli.converters` | `normalize_symbol`, `ensure_utc`, `parse_date_range`, `granularity_name` |
|
||||
| `mt5cli.exceptions` | `normalize_mt5_exception`, `call_with_normalized_errors`, `is_recoverable_mt5_error` |
|
||||
| Module | Examples |
|
||||
| ------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
|
||||
| `mt5cli.grafana` | `ensure_grafana_schema`, `create_grafana_views`, `create_grafana_indexes`, `create_snapshot_tables`, `start_snapshot_run`, `insert_account_snapshot`, `record_snapshot_run` |
|
||||
| `mt5cli.history` | `resolve_rate_view_name`, `resolve_rate_tables`, `load_rate_data`, `build_rate_view_name` |
|
||||
| `mt5cli.sdk` | `copy_rates_from`, `copy_ticks_from`, `account_info`, `symbols`, `mt5_summary`, `latest_rates` |
|
||||
| `mt5cli.schemas` | `DataKind`, `normalize_dataframe`, `validate_schema`, `DEDUP_KEYS` |
|
||||
| `mt5cli.utils` | `Dataset`, `IfExists`, `detect_format`, `export_dataframe`, `export_dataframe_to_sqlite` |
|
||||
| `mt5cli.converters` | `normalize_symbol`, `ensure_utc`, `parse_date_range`, `granularity_name` |
|
||||
| `mt5cli.exceptions` | `normalize_mt5_exception`, `call_with_normalized_errors`, `is_recoverable_mt5_error` |
|
||||
|
||||
## CLI commands
|
||||
|
||||
@@ -155,6 +213,15 @@ The Typer application in `mt5cli.cli` exposes file-export commands documented in
|
||||
- Delegate to the same Python APIs described here; they are not duplicated
|
||||
business logic.
|
||||
|
||||
`grafana-schema` initializes Grafana views, indexes, and snapshot tables in the
|
||||
target SQLite database without connecting to MT5. It is idempotent and safe to
|
||||
run repeatedly.
|
||||
|
||||
`snapshot` appends one timestamped row per enabled data type
|
||||
(`--with-account`, `--with-positions`, `--with-orders`, `--with-terminal`) and
|
||||
never places orders or modifies trading state. Both commands require
|
||||
`-o/--output` to point at a `.db` / SQLite file.
|
||||
|
||||
`order-send` is the expert raw-request path; it requires `--yes` and a fully
|
||||
constructed request payload. `close-positions` is the safer high-level helper
|
||||
that closes open positions by `--symbol` or `--ticket` using
|
||||
|
||||
@@ -35,6 +35,8 @@ from .sdk import (
|
||||
resolve_account_specs,
|
||||
update_history,
|
||||
update_history_with_config,
|
||||
update_observability,
|
||||
update_observability_with_config,
|
||||
)
|
||||
from .trading import (
|
||||
ExecutionStatus,
|
||||
@@ -140,6 +142,8 @@ __all__ = [
|
||||
"resolve_account_specs",
|
||||
"update_history",
|
||||
"update_history_with_config",
|
||||
"update_observability",
|
||||
"update_observability_with_config",
|
||||
"update_sltp_for_open_positions",
|
||||
"update_trailing_stop_loss_for_open_positions",
|
||||
]
|
||||
|
||||
@@ -802,6 +802,100 @@ def collect_history(
|
||||
)
|
||||
|
||||
|
||||
@app.command(rich_help_panel="Collection")
|
||||
def grafana_schema(ctx: typer.Context) -> None:
|
||||
"""Create or refresh Grafana-ready views and indexes in a SQLite database.
|
||||
|
||||
Idempotent — safe to run repeatedly on the same database. Requires SQLite
|
||||
output. Does not connect to MetaTrader 5.
|
||||
|
||||
Raises:
|
||||
typer.BadParameter: If the output format is not SQLite3.
|
||||
"""
|
||||
import sqlite3 as _sqlite3 # noqa: PLC0415
|
||||
|
||||
from .grafana import create_snapshot_tables, ensure_grafana_schema # noqa: PLC0415
|
||||
|
||||
export_ctx = _get_export_context(ctx)
|
||||
if export_ctx.output_format != "sqlite3":
|
||||
msg = (
|
||||
"grafana-schema requires SQLite3 output."
|
||||
" Use a .db/.sqlite/.sqlite3 extension or --format sqlite3."
|
||||
)
|
||||
raise typer.BadParameter(msg)
|
||||
with _sqlite3.connect(export_ctx.output) as conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
create_snapshot_tables(conn)
|
||||
ensure_grafana_schema(conn)
|
||||
logger.info("Grafana schema applied to %s", export_ctx.output)
|
||||
|
||||
|
||||
@app.command(rich_help_panel="Collection")
|
||||
def snapshot(
|
||||
ctx: typer.Context,
|
||||
symbol: Annotated[
|
||||
list[str] | None,
|
||||
typer.Option(
|
||||
"--symbol",
|
||||
"-s",
|
||||
help="Symbol filter for positions/orders (repeat for multiple).",
|
||||
),
|
||||
] = None,
|
||||
with_account: Annotated[
|
||||
bool,
|
||||
typer.Option("--with-account/--no-account", help="Snapshot account info."),
|
||||
] = True,
|
||||
with_positions: Annotated[
|
||||
bool,
|
||||
typer.Option(
|
||||
"--with-positions/--no-positions", help="Snapshot open positions."
|
||||
),
|
||||
] = True,
|
||||
with_orders: Annotated[
|
||||
bool,
|
||||
typer.Option("--with-orders/--no-orders", help="Snapshot active orders."),
|
||||
] = True,
|
||||
with_terminal: Annotated[
|
||||
bool,
|
||||
typer.Option("--with-terminal/--no-terminal", help="Snapshot terminal info."),
|
||||
] = True,
|
||||
with_grafana_schema: Annotated[
|
||||
bool,
|
||||
typer.Option(
|
||||
"--with-grafana-schema/--no-grafana-schema",
|
||||
help="Ensure Grafana views and indexes exist.",
|
||||
),
|
||||
] = False,
|
||||
) -> None:
|
||||
"""Snapshot current account, position, order, and terminal state into SQLite.
|
||||
|
||||
Appends a timestamped snapshot row for each data type. Never places
|
||||
orders or modifies trading state.
|
||||
|
||||
Raises:
|
||||
typer.BadParameter: If the output format is not SQLite3.
|
||||
"""
|
||||
export_ctx = _get_export_context(ctx)
|
||||
if export_ctx.output_format != "sqlite3":
|
||||
msg = (
|
||||
"snapshot requires SQLite3 output."
|
||||
" Use a .db/.sqlite/.sqlite3 extension or --format sqlite3."
|
||||
)
|
||||
raise typer.BadParameter(msg)
|
||||
sdk.update_observability_with_config(
|
||||
output=export_ctx.output,
|
||||
config=export_ctx.config,
|
||||
symbols=list(symbol) if symbol else None,
|
||||
include_account=with_account,
|
||||
include_positions=with_positions,
|
||||
include_orders=with_orders,
|
||||
include_terminal=with_terminal,
|
||||
with_grafana_schema=with_grafana_schema,
|
||||
)
|
||||
logger.info("Snapshot written to %s", export_ctx.output)
|
||||
|
||||
|
||||
def main() -> None:
|
||||
"""Run the mt5cli CLI."""
|
||||
app()
|
||||
|
||||
@@ -62,6 +62,8 @@ STABLE_SDK_EXPORTS: frozenset[str] = frozenset({
|
||||
"resolve_account_specs",
|
||||
"update_history",
|
||||
"update_history_with_config",
|
||||
"update_observability",
|
||||
"update_observability_with_config",
|
||||
"update_sltp_for_open_positions",
|
||||
"update_trailing_stop_loss_for_open_positions",
|
||||
})
|
||||
|
||||
@@ -0,0 +1,616 @@
|
||||
"""Grafana-oriented SQLite views, indexes, and snapshot tables."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import datetime
|
||||
import logging
|
||||
import sqlite3
|
||||
from typing import cast
|
||||
|
||||
from .history import get_table_columns
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_TRADE_DEAL_TYPES_SQL = "(0, 1)"
|
||||
|
||||
_GRAFANA_VIEW_NAMES = (
|
||||
"grafana_rates",
|
||||
"grafana_ticks",
|
||||
"grafana_history_deals",
|
||||
"grafana_history_orders",
|
||||
"grafana_trade_deals",
|
||||
"grafana_cash_events",
|
||||
"grafana_realized_pnl",
|
||||
"grafana_symbol_pnl",
|
||||
"grafana_trade_stats",
|
||||
"grafana_account_snapshots",
|
||||
"grafana_position_snapshots",
|
||||
"grafana_order_snapshots",
|
||||
"grafana_terminal_snapshots",
|
||||
)
|
||||
|
||||
|
||||
def _to_epoch_int(value: object) -> int | None:
|
||||
if value is None:
|
||||
return None
|
||||
if isinstance(value, datetime.datetime):
|
||||
return int(value.timestamp())
|
||||
if isinstance(value, (int, float)):
|
||||
return int(value)
|
||||
return None
|
||||
|
||||
|
||||
def _time_col_expr(col: str) -> str:
|
||||
return (
|
||||
f"CASE WHEN typeof(\"{col}\") IN ('integer', 'real')"
|
||||
f' THEN CAST("{col}" AS INTEGER)'
|
||||
f" ELSE CAST(strftime('%s', \"{col}\") AS INTEGER) END"
|
||||
)
|
||||
|
||||
|
||||
def _create_view_safe(
|
||||
conn: sqlite3.Connection,
|
||||
name: str,
|
||||
select_sql: str,
|
||||
) -> None:
|
||||
try:
|
||||
conn.execute(f'DROP VIEW IF EXISTS "{name}"')
|
||||
conn.execute(f'CREATE VIEW "{name}" AS {select_sql}')
|
||||
except sqlite3.Error as exc:
|
||||
logger.warning("Skipping view %s: %s", name, exc)
|
||||
|
||||
|
||||
def _other_cols(all_cols: set[str], exclude: set[str]) -> list[str]:
|
||||
return sorted(all_cols - exclude)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Snapshot table DDL
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
_SNAPSHOT_TABLE_DDLS: list[str] = [
|
||||
"""CREATE TABLE IF NOT EXISTS snapshot_runs (
|
||||
run_id INTEGER PRIMARY KEY,
|
||||
observed_at INTEGER NOT NULL,
|
||||
status TEXT NOT NULL,
|
||||
detail TEXT
|
||||
)""",
|
||||
"""CREATE TABLE IF NOT EXISTS account_snapshots (
|
||||
run_id INTEGER NOT NULL,
|
||||
login INTEGER,
|
||||
currency TEXT,
|
||||
balance REAL,
|
||||
equity REAL,
|
||||
margin REAL,
|
||||
margin_free REAL,
|
||||
margin_level REAL,
|
||||
profit REAL,
|
||||
leverage INTEGER
|
||||
)""",
|
||||
"""CREATE TABLE IF NOT EXISTS position_snapshots (
|
||||
run_id INTEGER NOT NULL,
|
||||
login INTEGER,
|
||||
ticket INTEGER,
|
||||
position_id INTEGER,
|
||||
symbol TEXT,
|
||||
type INTEGER,
|
||||
volume REAL,
|
||||
price_open REAL,
|
||||
price_current REAL,
|
||||
profit REAL,
|
||||
swap REAL,
|
||||
comment TEXT,
|
||||
magic INTEGER
|
||||
)""",
|
||||
"""CREATE TABLE IF NOT EXISTS order_snapshots (
|
||||
run_id INTEGER NOT NULL,
|
||||
login INTEGER,
|
||||
ticket INTEGER,
|
||||
symbol TEXT,
|
||||
type INTEGER,
|
||||
volume_current REAL,
|
||||
price_open REAL,
|
||||
price_current REAL,
|
||||
state INTEGER,
|
||||
comment TEXT,
|
||||
magic INTEGER,
|
||||
time_setup INTEGER
|
||||
)""",
|
||||
"""CREATE TABLE IF NOT EXISTS terminal_snapshots (
|
||||
run_id INTEGER NOT NULL,
|
||||
name TEXT,
|
||||
connected INTEGER,
|
||||
community_account INTEGER,
|
||||
trade_allowed INTEGER,
|
||||
trade_expert INTEGER,
|
||||
path TEXT,
|
||||
company TEXT,
|
||||
language TEXT
|
||||
)""",
|
||||
]
|
||||
|
||||
|
||||
def create_snapshot_tables(conn: sqlite3.Connection) -> None:
|
||||
"""Create snapshot tables idempotently."""
|
||||
for ddl in _SNAPSHOT_TABLE_DDLS:
|
||||
conn.execute(ddl)
|
||||
|
||||
|
||||
def start_snapshot_run(conn: sqlite3.Connection, observed_at: int) -> int:
|
||||
"""Insert a snapshot_runs row with status 'running' and return its run_id.
|
||||
|
||||
Returns:
|
||||
The auto-assigned run_id for the new row.
|
||||
"""
|
||||
cursor = conn.execute(
|
||||
"INSERT INTO snapshot_runs (observed_at, status) VALUES (?, 'running')",
|
||||
(observed_at,),
|
||||
)
|
||||
return cast("int", cursor.lastrowid)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# View builders
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def _build_grafana_rates(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "rates")
|
||||
required = {"time", "symbol", "timeframe"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_rates: rates table missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
time_expr = _time_col_expr("time")
|
||||
others = _other_cols(cols, {"time"})
|
||||
other_sql = ", ".join(f'"{c}"' for c in others)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_rates",
|
||||
f'SELECT {time_expr} AS "time", {other_sql} FROM "rates"', # noqa: S608
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_ticks(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "ticks")
|
||||
required = {"time", "symbol"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_ticks: ticks table missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
time_expr = _time_col_expr("time")
|
||||
others = _other_cols(cols, {"time"})
|
||||
other_sql = ", ".join(f'"{c}"' for c in others)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_ticks",
|
||||
f'SELECT {time_expr} AS "time", {other_sql} FROM "ticks"', # noqa: S608
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_history_deals(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_deals")
|
||||
if "time" not in cols:
|
||||
logger.warning("Skipping grafana_history_deals: history_deals.time is missing")
|
||||
return
|
||||
time_expr = _time_col_expr("time")
|
||||
others = _other_cols(cols, {"time"})
|
||||
other_sql = ", ".join(f'"{c}"' for c in others)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_history_deals",
|
||||
f'SELECT {time_expr} AS "time", {other_sql} FROM "history_deals"', # noqa: S608
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_history_orders(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_orders")
|
||||
if "time_setup" not in cols:
|
||||
logger.warning(
|
||||
"Skipping grafana_history_orders: history_orders.time_setup is missing"
|
||||
)
|
||||
return
|
||||
time_expr = _time_col_expr("time_setup")
|
||||
others = _other_cols(cols, set())
|
||||
other_sql = ", ".join(f'"{c}"' for c in others)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_history_orders",
|
||||
f'SELECT {time_expr} AS "time", {other_sql} FROM "history_orders"', # noqa: S608
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_trade_deals(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_deals")
|
||||
required = {"time", "type"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_trade_deals: history_deals missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
time_expr = _time_col_expr("time")
|
||||
others = _other_cols(cols, {"time"})
|
||||
other_sql = ", ".join(f'"{c}"' for c in others)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_trade_deals",
|
||||
f'SELECT {time_expr} AS "time", {other_sql}' # noqa: S608
|
||||
f' FROM "history_deals" WHERE "type" IN {_TRADE_DEAL_TYPES_SQL}',
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_cash_events(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_deals")
|
||||
required = {"time", "type"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_cash_events: history_deals missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
time_expr = _time_col_expr("time")
|
||||
others = _other_cols(cols, {"time"})
|
||||
other_sql = ", ".join(f'"{c}"' for c in others)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_cash_events",
|
||||
f'SELECT {time_expr} AS "time", {other_sql}' # noqa: S608
|
||||
f' FROM "history_deals" WHERE "type" NOT IN {_TRADE_DEAL_TYPES_SQL}',
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_realized_pnl(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_deals")
|
||||
required = {"symbol", "profit", "type", "entry"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_realized_pnl: history_deals missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_realized_pnl",
|
||||
'SELECT "symbol",' # noqa: S608
|
||||
' SUM("profit") AS cumulative_pnl, COUNT(*) AS deal_count'
|
||||
' FROM "history_deals"'
|
||||
f' WHERE "type" IN {_TRADE_DEAL_TYPES_SQL}'
|
||||
' AND "entry" IN (1, 2, 3)'
|
||||
' AND "symbol" IS NOT NULL AND "symbol" != \'\''
|
||||
' GROUP BY "symbol"',
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_symbol_pnl(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_deals")
|
||||
required = {"time", "symbol", "profit", "type", "entry"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_symbol_pnl: history_deals missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
time_expr = _time_col_expr("time")
|
||||
select_parts = [f'{time_expr} AS "time"', '"symbol"', '"profit"']
|
||||
if "volume" in cols:
|
||||
select_parts.append('"volume"')
|
||||
if "price" in cols:
|
||||
select_parts.append('"price"')
|
||||
select_sql = ", ".join(select_parts)
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_symbol_pnl",
|
||||
f'SELECT {select_sql} FROM "history_deals"' # noqa: S608
|
||||
f' WHERE "type" IN {_TRADE_DEAL_TYPES_SQL}'
|
||||
' AND "entry" IN (1, 2, 3)'
|
||||
' AND "symbol" IS NOT NULL AND "symbol" != \'\'',
|
||||
)
|
||||
|
||||
|
||||
def _build_grafana_trade_stats(conn: sqlite3.Connection) -> None:
|
||||
cols = get_table_columns(conn, "history_deals")
|
||||
required = {"symbol", "profit", "type"}
|
||||
if not required.issubset(cols):
|
||||
logger.warning(
|
||||
"Skipping grafana_trade_stats: history_deals missing columns %s",
|
||||
sorted(required - cols),
|
||||
)
|
||||
return
|
||||
has_entry = "entry" in cols
|
||||
entry_filter = ' AND "entry" IN (1, 2, 3)' if has_entry else ""
|
||||
_create_view_safe(
|
||||
conn,
|
||||
"grafana_trade_stats",
|
||||
'SELECT "symbol",' # noqa: S608
|
||||
" COUNT(*) AS total_deals,"
|
||||
' SUM(CASE WHEN "profit" > 0 THEN 1 ELSE 0 END) AS winning_deals,'
|
||||
' SUM(CASE WHEN "profit" <= 0 THEN 1 ELSE 0 END) AS losing_deals,'
|
||||
' SUM("profit") AS total_profit,'
|
||||
' AVG("profit") AS avg_profit,'
|
||||
' MAX("profit") AS max_profit,'
|
||||
' MIN("profit") AS min_profit'
|
||||
' FROM "history_deals"'
|
||||
f' WHERE "type" IN {_TRADE_DEAL_TYPES_SQL}'
|
||||
f"{entry_filter}"
|
||||
' AND "symbol" IS NOT NULL AND "symbol" != \'\''
|
||||
' GROUP BY "symbol"',
|
||||
)
|
||||
|
||||
|
||||
def _build_snapshot_view(
|
||||
conn: sqlite3.Connection,
|
||||
view_name: str,
|
||||
table_name: str,
|
||||
) -> None:
|
||||
cols = get_table_columns(conn, table_name)
|
||||
if not cols:
|
||||
logger.warning("Skipping %s: %s table missing", view_name, table_name)
|
||||
return
|
||||
if "run_id" not in cols:
|
||||
logger.warning("Skipping %s: %s missing run_id column", view_name, table_name)
|
||||
return
|
||||
others = _other_cols(cols, {"run_id"})
|
||||
run_cols = get_table_columns(conn, "snapshot_runs")
|
||||
if {"run_id", "observed_at", "status"}.issubset(run_cols):
|
||||
other_sql = (", " + ", ".join(f's."{c}"' for c in others)) if others else ""
|
||||
select_cols = f'r."observed_at" AS "time", s."run_id"{other_sql}'
|
||||
_create_view_safe(
|
||||
conn,
|
||||
view_name,
|
||||
f'SELECT {select_cols} FROM "{table_name}" s' # noqa: S608
|
||||
f' JOIN "snapshot_runs" r ON s."run_id" = r."run_id"'
|
||||
f" WHERE r.\"status\" = 'ok'",
|
||||
)
|
||||
else:
|
||||
logger.warning("Skipping %s: snapshot_runs missing required columns", view_name)
|
||||
|
||||
|
||||
def _build_grafana_account_snapshots(conn: sqlite3.Connection) -> None:
|
||||
_build_snapshot_view(conn, "grafana_account_snapshots", "account_snapshots")
|
||||
|
||||
|
||||
def _build_grafana_position_snapshots(conn: sqlite3.Connection) -> None:
|
||||
_build_snapshot_view(conn, "grafana_position_snapshots", "position_snapshots")
|
||||
|
||||
|
||||
def _build_grafana_order_snapshots(conn: sqlite3.Connection) -> None:
|
||||
_build_snapshot_view(conn, "grafana_order_snapshots", "order_snapshots")
|
||||
|
||||
|
||||
def _build_grafana_terminal_snapshots(conn: sqlite3.Connection) -> None:
|
||||
_build_snapshot_view(conn, "grafana_terminal_snapshots", "terminal_snapshots")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Public API
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def create_grafana_views(conn: sqlite3.Connection) -> None:
|
||||
"""Create all Grafana-facing views idempotently.
|
||||
|
||||
Missing source tables cause the affected view to be skipped with a warning;
|
||||
other views are unaffected. Stale views whose source table or required
|
||||
columns have disappeared are dropped before rebuild.
|
||||
"""
|
||||
for name in _GRAFANA_VIEW_NAMES:
|
||||
conn.execute(f'DROP VIEW IF EXISTS "{name}"')
|
||||
_build_grafana_rates(conn)
|
||||
_build_grafana_ticks(conn)
|
||||
_build_grafana_history_deals(conn)
|
||||
_build_grafana_history_orders(conn)
|
||||
_build_grafana_trade_deals(conn)
|
||||
_build_grafana_cash_events(conn)
|
||||
_build_grafana_realized_pnl(conn)
|
||||
_build_grafana_symbol_pnl(conn)
|
||||
_build_grafana_trade_stats(conn)
|
||||
_build_grafana_account_snapshots(conn)
|
||||
_build_grafana_position_snapshots(conn)
|
||||
_build_grafana_order_snapshots(conn)
|
||||
_build_grafana_terminal_snapshots(conn)
|
||||
|
||||
|
||||
def create_grafana_indexes(conn: sqlite3.Connection) -> None:
|
||||
"""Create Grafana query performance indexes idempotently."""
|
||||
rates_cols = get_table_columns(conn, "rates")
|
||||
if {"time", "symbol", "timeframe"}.issubset(rates_cols):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_rates_time_symbol_timeframe"
|
||||
' ON "rates"("time", "symbol", "timeframe")',
|
||||
)
|
||||
|
||||
ticks_cols = get_table_columns(conn, "ticks")
|
||||
if {"time", "symbol"}.issubset(ticks_cols):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_ticks_time_symbol"
|
||||
' ON "ticks"("time", "symbol")',
|
||||
)
|
||||
|
||||
deals_cols = get_table_columns(conn, "history_deals")
|
||||
if {"time", "symbol"}.issubset(deals_cols):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_history_deals_time_symbol"
|
||||
' ON "history_deals"("time", "symbol")',
|
||||
)
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_history_deals_symbol_time"
|
||||
' ON "history_deals"("symbol", "time")',
|
||||
)
|
||||
|
||||
orders_cols = get_table_columns(conn, "history_orders")
|
||||
if {"time_setup", "symbol"}.issubset(orders_cols):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_history_orders_time_setup_symbol"
|
||||
' ON "history_orders"("time_setup", "symbol")',
|
||||
)
|
||||
|
||||
if {"run_id", "login"}.issubset(get_table_columns(conn, "account_snapshots")):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_account_snapshots_time_login"
|
||||
' ON "account_snapshots"("run_id", "login")',
|
||||
)
|
||||
if {"run_id", "symbol"}.issubset(get_table_columns(conn, "position_snapshots")):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_position_snapshots_time_symbol"
|
||||
' ON "position_snapshots"("run_id", "symbol")',
|
||||
)
|
||||
if {"run_id", "symbol"}.issubset(get_table_columns(conn, "order_snapshots")):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_order_snapshots_time_symbol"
|
||||
' ON "order_snapshots"("run_id", "symbol")',
|
||||
)
|
||||
if {"observed_at", "status"}.issubset(get_table_columns(conn, "snapshot_runs")):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_snapshot_runs_time_status"
|
||||
' ON "snapshot_runs"("observed_at", "status")',
|
||||
)
|
||||
|
||||
|
||||
def ensure_grafana_schema(conn: sqlite3.Connection) -> None:
|
||||
"""Create snapshot tables, Grafana views, and indexes idempotently."""
|
||||
create_snapshot_tables(conn)
|
||||
create_grafana_views(conn)
|
||||
create_grafana_indexes(conn)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Snapshot insert helpers
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def insert_account_snapshot(
|
||||
conn: sqlite3.Connection,
|
||||
run_id: int,
|
||||
row: dict[str, object],
|
||||
) -> None:
|
||||
"""Append one account state row to account_snapshots."""
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots"
|
||||
" (run_id, login, currency, balance, equity,"
|
||||
" margin, margin_free, margin_level, profit, leverage)"
|
||||
" VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
(
|
||||
run_id,
|
||||
row.get("login"),
|
||||
row.get("currency"),
|
||||
row.get("balance"),
|
||||
row.get("equity"),
|
||||
row.get("margin"),
|
||||
row.get("margin_free"),
|
||||
row.get("margin_level"),
|
||||
row.get("profit"),
|
||||
row.get("leverage"),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def insert_position_snapshots(
|
||||
conn: sqlite3.Connection,
|
||||
run_id: int,
|
||||
login: int | None,
|
||||
rows: list[dict[str, object]],
|
||||
) -> None:
|
||||
"""Append position rows to position_snapshots; no-op when rows is empty."""
|
||||
if not rows:
|
||||
return
|
||||
conn.executemany(
|
||||
"INSERT INTO position_snapshots"
|
||||
" (run_id, login, ticket, position_id, symbol, type, volume,"
|
||||
" price_open, price_current, profit, swap, comment, magic)"
|
||||
" VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
[
|
||||
(
|
||||
run_id,
|
||||
login,
|
||||
r.get("ticket"),
|
||||
r.get("position_id"),
|
||||
r.get("symbol"),
|
||||
r.get("type"),
|
||||
r.get("volume"),
|
||||
r.get("price_open"),
|
||||
r.get("price_current"),
|
||||
r.get("profit"),
|
||||
r.get("swap"),
|
||||
r.get("comment"),
|
||||
r.get("magic"),
|
||||
)
|
||||
for r in rows
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
def insert_order_snapshots(
|
||||
conn: sqlite3.Connection,
|
||||
run_id: int,
|
||||
login: int | None,
|
||||
rows: list[dict[str, object]],
|
||||
) -> None:
|
||||
"""Append order rows to order_snapshots; no-op when rows is empty."""
|
||||
if not rows:
|
||||
return
|
||||
conn.executemany(
|
||||
"INSERT INTO order_snapshots"
|
||||
" (run_id, login, ticket, symbol, type, volume_current,"
|
||||
" price_open, price_current, state, comment, magic, time_setup)"
|
||||
" VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
[
|
||||
(
|
||||
run_id,
|
||||
login,
|
||||
r.get("ticket"),
|
||||
r.get("symbol"),
|
||||
r.get("type"),
|
||||
r.get("volume_current"),
|
||||
r.get("price_open"),
|
||||
r.get("price_current"),
|
||||
r.get("state"),
|
||||
r.get("comment"),
|
||||
r.get("magic"),
|
||||
_to_epoch_int(r.get("time_setup")),
|
||||
)
|
||||
for r in rows
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
def insert_terminal_snapshot(
|
||||
conn: sqlite3.Connection,
|
||||
run_id: int,
|
||||
row: dict[str, object],
|
||||
) -> None:
|
||||
"""Append one terminal state row to terminal_snapshots."""
|
||||
conn.execute(
|
||||
"INSERT INTO terminal_snapshots"
|
||||
" (run_id, name, connected, community_account,"
|
||||
" trade_allowed, trade_expert, path, company, language)"
|
||||
" VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
(
|
||||
run_id,
|
||||
row.get("name"),
|
||||
row.get("connected"),
|
||||
row.get("community_account"),
|
||||
row.get("trade_allowed"),
|
||||
row.get("trade_expert"),
|
||||
row.get("path"),
|
||||
row.get("company"),
|
||||
row.get("language"),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def record_snapshot_run(
|
||||
conn: sqlite3.Connection,
|
||||
run_id: int,
|
||||
status: str,
|
||||
detail: str | None = None,
|
||||
) -> None:
|
||||
"""Finalize a snapshot run by setting its status."""
|
||||
conn.execute(
|
||||
"UPDATE snapshot_runs SET status = ?, detail = ? WHERE run_id = ?",
|
||||
(status, detail, run_id),
|
||||
)
|
||||
+181
-3
@@ -8,7 +8,7 @@ import os
|
||||
import re
|
||||
import sqlite3
|
||||
import time
|
||||
from contextlib import contextmanager
|
||||
from contextlib import closing, contextmanager
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
@@ -22,6 +22,16 @@ try:
|
||||
except ImportError: # pragma: no cover
|
||||
Mt5TradingError = None # type: ignore[assignment]
|
||||
|
||||
from .grafana import (
|
||||
create_snapshot_tables,
|
||||
ensure_grafana_schema,
|
||||
insert_account_snapshot,
|
||||
insert_order_snapshots,
|
||||
insert_position_snapshots,
|
||||
insert_terminal_snapshot,
|
||||
record_snapshot_run,
|
||||
start_snapshot_run,
|
||||
)
|
||||
from .history import (
|
||||
create_cash_events_view,
|
||||
create_history_indexes,
|
||||
@@ -154,6 +164,8 @@ __all__ = [
|
||||
"terminal_info",
|
||||
"update_history",
|
||||
"update_history_with_config",
|
||||
"update_observability",
|
||||
"update_observability_with_config",
|
||||
"version",
|
||||
]
|
||||
|
||||
@@ -1011,7 +1023,7 @@ def update_history( # noqa: PLR0913
|
||||
sorted(dataset.value for dataset in request.selected),
|
||||
request.output_path,
|
||||
)
|
||||
with sqlite3.connect(request.output_path) as conn:
|
||||
with closing(sqlite3.connect(request.output_path)) as conn, conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
write_incremental_datasets(
|
||||
@@ -1258,7 +1270,11 @@ def collect_history(
|
||||
tf = _coerce_timeframe(timeframe)
|
||||
tick_flags = _coerce_tick_flags(flags)
|
||||
mt5_config = config or build_config()
|
||||
with connected_client(mt5_config) as client, sqlite3.connect(output) as conn:
|
||||
with (
|
||||
connected_client(mt5_config) as client,
|
||||
closing(sqlite3.connect(output)) as conn,
|
||||
conn,
|
||||
):
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
written_tables, written_columns = write_collected_datasets(
|
||||
@@ -2148,3 +2164,165 @@ def mt5_summary(*, config: Mt5Config | None = None) -> dict[str, object]:
|
||||
def mt5_summary_as_df(*, config: Mt5Config | None = None) -> pd.DataFrame:
|
||||
"""Return an export-safe terminal/account status summary DataFrame."""
|
||||
return _make_client(config=config).mt5_summary_as_df()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Observability: account / position / order / terminal snapshots
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def _snapshot_account(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
run_id: int,
|
||||
) -> int | None:
|
||||
df = client.account_info_as_df()
|
||||
if df.empty:
|
||||
logger.warning(
|
||||
"account_info_as_df returned empty frame; skipping account snapshot"
|
||||
)
|
||||
return None
|
||||
row = cast("dict[str, object]", df.iloc[0].to_dict())
|
||||
insert_account_snapshot(conn, run_id, row)
|
||||
login_val = row.get("login")
|
||||
return int(login_val) if login_val is not None else None # type: ignore[arg-type]
|
||||
|
||||
|
||||
def _snapshot_positions(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
run_id: int,
|
||||
login: int | None,
|
||||
symbols: Sequence[str] | None,
|
||||
) -> None:
|
||||
df: pd.DataFrame = client.positions_get_as_df()
|
||||
if symbols is not None and not df.empty and "symbol" in df.columns:
|
||||
df = df[df["symbol"].isin(symbols)].reset_index(drop=True)
|
||||
raw = df.to_dict(orient="records") if not df.empty else []
|
||||
rows = cast("list[dict[str, object]]", raw)
|
||||
insert_position_snapshots(conn, run_id, login, rows)
|
||||
|
||||
|
||||
def _snapshot_orders(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
run_id: int,
|
||||
login: int | None,
|
||||
symbols: Sequence[str] | None,
|
||||
) -> None:
|
||||
df: pd.DataFrame = client.orders_get_as_df()
|
||||
if symbols is not None and not df.empty and "symbol" in df.columns:
|
||||
df = df[df["symbol"].isin(symbols)].reset_index(drop=True)
|
||||
raw = df.to_dict(orient="records") if not df.empty else []
|
||||
rows = cast("list[dict[str, object]]", raw)
|
||||
insert_order_snapshots(conn, run_id, login, rows)
|
||||
|
||||
|
||||
def _snapshot_terminal(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
run_id: int,
|
||||
) -> None:
|
||||
df = client.terminal_info_as_df()
|
||||
if df.empty:
|
||||
logger.warning(
|
||||
"terminal_info_as_df returned empty frame; skipping terminal snapshot"
|
||||
)
|
||||
return
|
||||
row = cast("dict[str, object]", df.iloc[0].to_dict())
|
||||
insert_terminal_snapshot(conn, run_id, row)
|
||||
|
||||
|
||||
def update_observability(
|
||||
*,
|
||||
client: Mt5DataClient,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str] | None = None,
|
||||
include_account: bool = True,
|
||||
include_positions: bool = True,
|
||||
include_orders: bool = True,
|
||||
include_terminal: bool = True,
|
||||
with_grafana_schema: bool = False,
|
||||
) -> None:
|
||||
"""Snapshot current account/position/order/terminal state into SQLite.
|
||||
|
||||
Reads the current MT5 state and appends timestamped snapshot rows. Never
|
||||
places orders or modifies trading state.
|
||||
|
||||
Args:
|
||||
client: Connected MT5 data client.
|
||||
output: SQLite database path.
|
||||
symbols: Optional symbol filter for positions and orders. When None,
|
||||
all positions and orders are snapshotted.
|
||||
include_account: Snapshot account info into ``account_snapshots``.
|
||||
include_positions: Snapshot open positions into ``position_snapshots``.
|
||||
include_orders: Snapshot active orders into ``order_snapshots``.
|
||||
include_terminal: Snapshot terminal info into ``terminal_snapshots``.
|
||||
with_grafana_schema: Ensure Grafana views and indexes exist. Defaults
|
||||
to ``False``; run ``grafana-schema`` once to set up the schema,
|
||||
then use ``snapshot`` repeatedly without this flag.
|
||||
"""
|
||||
observed_at = int(datetime.now(UTC).timestamp())
|
||||
with closing(sqlite3.connect(Path(output))) as conn, conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
if with_grafana_schema:
|
||||
ensure_grafana_schema(conn)
|
||||
else:
|
||||
create_snapshot_tables(conn)
|
||||
run_id = start_snapshot_run(conn, observed_at)
|
||||
login: int | None = None
|
||||
try:
|
||||
if include_account:
|
||||
login = _snapshot_account(conn, client, run_id)
|
||||
if include_positions:
|
||||
_snapshot_positions(conn, client, run_id, login, symbols)
|
||||
if include_orders:
|
||||
_snapshot_orders(conn, client, run_id, login, symbols)
|
||||
if include_terminal:
|
||||
_snapshot_terminal(conn, client, run_id)
|
||||
record_snapshot_run(conn, run_id, "ok")
|
||||
except Exception:
|
||||
record_snapshot_run(conn, run_id, "error")
|
||||
conn.commit()
|
||||
raise
|
||||
|
||||
|
||||
def update_observability_with_config(
|
||||
*,
|
||||
output: Path | str,
|
||||
config: Mt5Config | None = None,
|
||||
symbols: Sequence[str] | None = None,
|
||||
include_account: bool = True,
|
||||
include_positions: bool = True,
|
||||
include_orders: bool = True,
|
||||
include_terminal: bool = True,
|
||||
with_grafana_schema: bool = False,
|
||||
) -> None:
|
||||
"""Snapshot current MT5 state, opening and closing the MT5 connection.
|
||||
|
||||
Convenience wrapper around :func:`update_observability` for standalone use.
|
||||
|
||||
Args:
|
||||
output: SQLite database path.
|
||||
config: MT5 connection configuration. Defaults to an empty config that
|
||||
attaches to a running terminal.
|
||||
symbols: Optional symbol filter for positions and orders.
|
||||
include_account: Snapshot account info.
|
||||
include_positions: Snapshot open positions.
|
||||
include_orders: Snapshot active orders.
|
||||
include_terminal: Snapshot terminal info.
|
||||
with_grafana_schema: Ensure Grafana views and indexes exist.
|
||||
"""
|
||||
mt5_config = config or build_config()
|
||||
with connected_client(mt5_config) as client:
|
||||
update_observability(
|
||||
client=client,
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
include_account=include_account,
|
||||
include_positions=include_positions,
|
||||
include_orders=include_orders,
|
||||
include_terminal=include_terminal,
|
||||
with_grafana_schema=with_grafana_schema,
|
||||
)
|
||||
|
||||
+2
-1
@@ -4,6 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import json
|
||||
import sqlite3
|
||||
from contextlib import closing
|
||||
from datetime import UTC, datetime
|
||||
from enum import StrEnum
|
||||
from pathlib import Path
|
||||
@@ -277,7 +278,7 @@ def export_dataframe_to_sqlite(
|
||||
full table, so repeated appends cost O(table size); index the key
|
||||
columns when appending frequently.
|
||||
"""
|
||||
with sqlite3.connect(output_path) as conn:
|
||||
with closing(sqlite3.connect(output_path)) as conn, conn:
|
||||
df.to_sql( # type: ignore[reportUnknownMemberType]
|
||||
table_name,
|
||||
conn,
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "mt5cli"
|
||||
version = "1.0.3"
|
||||
version = "1.1.0"
|
||||
description = "Generic MT5 data and execution infrastructure for Python applications"
|
||||
authors = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||
maintainers = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||
|
||||
@@ -2,12 +2,17 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
from typing import TYPE_CHECKING, Any, Literal
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pandas as pd
|
||||
import pytest
|
||||
from pytest_mock import MockerFixture # noqa: TC002
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from types import TracebackType
|
||||
|
||||
_DATAFRAME_METHODS = (
|
||||
"copy_rates_from_as_df",
|
||||
"copy_rates_from_pos_as_df",
|
||||
@@ -30,6 +35,25 @@ _DATAFRAME_METHODS = (
|
||||
"order_send_as_df",
|
||||
)
|
||||
|
||||
_ORIGINAL_SQLITE_CONNECT = sqlite3.connect
|
||||
|
||||
|
||||
class ClosingSqliteConnection(sqlite3.Connection):
|
||||
"""SQLite connection that closes after context-manager exit in tests."""
|
||||
|
||||
def __exit__(
|
||||
self,
|
||||
exc_type: type[BaseException] | None,
|
||||
exc_value: BaseException | None,
|
||||
traceback: TracebackType | None,
|
||||
) -> Literal[False]:
|
||||
"""Commit or roll back the transaction, then close the connection."""
|
||||
try:
|
||||
super().__exit__(exc_type, exc_value, traceback)
|
||||
finally:
|
||||
self.close()
|
||||
return False
|
||||
|
||||
|
||||
def build_mock_mt5_data_client() -> MagicMock:
|
||||
"""Return a MagicMock Mt5DataClient with common DataFrame stubs."""
|
||||
@@ -50,3 +74,17 @@ def mock_client(mocker: MockerFixture) -> MagicMock:
|
||||
client = build_mock_mt5_data_client()
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=client)
|
||||
return client
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def close_sqlite_context_connections(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Make test SQLite context managers close their connection handles."""
|
||||
|
||||
def connect(
|
||||
*args: Any, # noqa: ANN401
|
||||
**kwargs: Any, # noqa: ANN401
|
||||
) -> sqlite3.Connection:
|
||||
kwargs.setdefault("factory", ClosingSqliteConnection)
|
||||
return _ORIGINAL_SQLITE_CONNECT(*args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(sqlite3, "connect", connect)
|
||||
|
||||
@@ -1855,6 +1855,174 @@ class TestCollectHistory:
|
||||
)
|
||||
|
||||
|
||||
class TestGrafanaSchemaCommand:
|
||||
"""Tests for the grafana-schema CLI command."""
|
||||
|
||||
def test_grafana_schema_creates_snapshot_tables_in_sqlite(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""grafana-schema applies Grafana schema to a SQLite database."""
|
||||
output = tmp_path / "out.db"
|
||||
result = runner.invoke(app, ["-o", str(output), "grafana-schema"])
|
||||
assert result.exit_code == 0, result.output
|
||||
with sqlite3.connect(output) as conn:
|
||||
tables = {
|
||||
row[0]
|
||||
for row in conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type='table'"
|
||||
).fetchall()
|
||||
}
|
||||
assert "snapshot_runs" in tables
|
||||
assert "account_snapshots" in tables
|
||||
|
||||
def test_grafana_schema_is_idempotent(self, tmp_path: Path) -> None:
|
||||
"""grafana-schema can be invoked multiple times without error."""
|
||||
output = tmp_path / "out.db"
|
||||
result1 = runner.invoke(app, ["-o", str(output), "grafana-schema"])
|
||||
result2 = runner.invoke(app, ["-o", str(output), "grafana-schema"])
|
||||
assert result1.exit_code == 0, result1.output
|
||||
assert result2.exit_code == 0, result2.output
|
||||
|
||||
def test_grafana_schema_rejects_non_sqlite_output(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""grafana-schema fails when output is not a SQLite3 format."""
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.csv"), "grafana-schema"],
|
||||
)
|
||||
assert result.exit_code != 0
|
||||
assert "grafana-schema requires SQLite3 output" in result.output
|
||||
|
||||
|
||||
class TestSnapshotCommand:
|
||||
"""Tests for the snapshot CLI command."""
|
||||
|
||||
def test_snapshot_rejects_non_sqlite_output(self, tmp_path: Path) -> None:
|
||||
"""Snapshot fails when output is not a SQLite3 format."""
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.csv"), "snapshot"],
|
||||
)
|
||||
assert result.exit_code != 0
|
||||
assert "snapshot requires SQLite3 output" in result.output
|
||||
|
||||
def test_snapshot_delegates_to_update_observability_with_config(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Snapshot calls sdk.update_observability_with_config."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
output = tmp_path / "out.db"
|
||||
result = runner.invoke(app, ["-o", str(output), "snapshot"])
|
||||
assert result.exit_code == 0, result.output
|
||||
updater.assert_called_once()
|
||||
kwargs = updater.call_args.kwargs
|
||||
assert kwargs["output"] == output
|
||||
assert kwargs["symbols"] is None
|
||||
assert kwargs["include_account"] is True
|
||||
assert kwargs["include_positions"] is True
|
||||
assert kwargs["include_orders"] is True
|
||||
assert kwargs["include_terminal"] is True
|
||||
assert kwargs["with_grafana_schema"] is False
|
||||
|
||||
def test_snapshot_with_symbol_filter(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""Snapshot passes symbol list to update_observability_with_config."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
result = runner.invoke(
|
||||
app,
|
||||
[
|
||||
"-o",
|
||||
str(tmp_path / "out.db"),
|
||||
"snapshot",
|
||||
"--symbol",
|
||||
"EURUSD",
|
||||
"--symbol",
|
||||
"GBPUSD",
|
||||
],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
kwargs = updater.call_args.kwargs
|
||||
assert kwargs["symbols"] == ["EURUSD", "GBPUSD"]
|
||||
|
||||
def test_snapshot_with_no_account_flag(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""--no-account disables account snapshotting."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.db"), "snapshot", "--no-account"],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
assert updater.call_args.kwargs["include_account"] is False
|
||||
|
||||
def test_snapshot_with_no_positions_flag(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""--no-positions disables position snapshotting."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.db"), "snapshot", "--no-positions"],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
assert updater.call_args.kwargs["include_positions"] is False
|
||||
|
||||
def test_snapshot_with_no_orders_flag(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""--no-orders disables order snapshotting."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.db"), "snapshot", "--no-orders"],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
assert updater.call_args.kwargs["include_orders"] is False
|
||||
|
||||
def test_snapshot_with_no_terminal_flag(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""--no-terminal disables terminal snapshotting."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.db"), "snapshot", "--no-terminal"],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
assert updater.call_args.kwargs["include_terminal"] is False
|
||||
|
||||
def test_snapshot_with_no_grafana_schema_flag(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
mocker: MockerFixture,
|
||||
) -> None:
|
||||
"""--no-grafana-schema disables Grafana schema creation."""
|
||||
updater = mocker.patch("mt5cli.cli.sdk.update_observability_with_config")
|
||||
result = runner.invoke(
|
||||
app,
|
||||
["-o", str(tmp_path / "out.db"), "snapshot", "--no-grafana-schema"],
|
||||
)
|
||||
assert result.exit_code == 0, result.output
|
||||
assert updater.call_args.kwargs["with_grafana_schema"] is False
|
||||
|
||||
|
||||
class TestMain:
|
||||
"""Tests for the main entry point."""
|
||||
|
||||
|
||||
@@ -0,0 +1,838 @@
|
||||
"""Tests for mt5cli.grafana module."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import sqlite3
|
||||
from typing import TYPE_CHECKING
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pandas as pd
|
||||
import pytest
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import Iterator
|
||||
|
||||
from mt5cli.grafana import (
|
||||
_build_snapshot_view, # type: ignore[reportPrivateUsage]
|
||||
_create_view_safe, # type: ignore[reportPrivateUsage]
|
||||
create_grafana_indexes,
|
||||
create_grafana_views,
|
||||
create_snapshot_tables,
|
||||
ensure_grafana_schema,
|
||||
insert_account_snapshot,
|
||||
insert_order_snapshots,
|
||||
insert_position_snapshots,
|
||||
insert_terminal_snapshot,
|
||||
record_snapshot_run,
|
||||
start_snapshot_run,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def conn() -> Iterator[sqlite3.Connection]:
|
||||
"""Yield an in-memory SQLite connection for each test."""
|
||||
with sqlite3.connect(":memory:") as c:
|
||||
yield c
|
||||
|
||||
|
||||
def _get_names(conn: sqlite3.Connection, type_: str) -> set[str]:
|
||||
return {
|
||||
row[0]
|
||||
for row in conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type=?",
|
||||
(type_,),
|
||||
).fetchall()
|
||||
}
|
||||
|
||||
|
||||
def _make_rates_table(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"CREATE TABLE rates"
|
||||
" (time TEXT, symbol TEXT, timeframe INTEGER,"
|
||||
" open REAL, high REAL, low REAL, close REAL)"
|
||||
)
|
||||
|
||||
|
||||
def _make_ticks_table(conn: sqlite3.Connection) -> None:
|
||||
conn.execute("CREATE TABLE ticks (time TEXT, symbol TEXT, bid REAL, ask REAL)")
|
||||
|
||||
|
||||
def _make_history_deals_full(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"CREATE TABLE history_deals"
|
||||
" (time TEXT, symbol TEXT, profit REAL, type INTEGER,"
|
||||
" entry INTEGER, volume REAL, price REAL, ticket INTEGER, position_id INTEGER)"
|
||||
)
|
||||
|
||||
|
||||
def _make_history_deals_minimal(conn: sqlite3.Connection) -> None:
|
||||
"""history_deals with only time, type, symbol, profit — no entry/volume/price."""
|
||||
conn.execute(
|
||||
"CREATE TABLE history_deals (time TEXT, symbol TEXT, profit REAL, type INTEGER)"
|
||||
)
|
||||
|
||||
|
||||
def _make_history_orders_table(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(
|
||||
"CREATE TABLE history_orders"
|
||||
" (time_setup TEXT, symbol TEXT, ticket INTEGER, type INTEGER)"
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestSnapshotTables
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestSnapshotTables:
|
||||
"""Tests for create_snapshot_tables."""
|
||||
|
||||
def test_creates_all_five_tables(self, conn: sqlite3.Connection) -> None:
|
||||
"""All five snapshot tables are created."""
|
||||
create_snapshot_tables(conn)
|
||||
tables = _get_names(conn, "table")
|
||||
assert "snapshot_runs" in tables
|
||||
assert "account_snapshots" in tables
|
||||
assert "position_snapshots" in tables
|
||||
assert "order_snapshots" in tables
|
||||
assert "terminal_snapshots" in tables
|
||||
|
||||
def test_is_idempotent(self, conn: sqlite3.Connection) -> None:
|
||||
"""Calling create_snapshot_tables twice does not raise."""
|
||||
create_snapshot_tables(conn)
|
||||
create_snapshot_tables(conn)
|
||||
tables = _get_names(conn, "table")
|
||||
assert "snapshot_runs" in tables
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestCreateViewSafe
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestCreateViewSafe:
|
||||
"""Tests for _create_view_safe."""
|
||||
|
||||
def test_creates_view_successfully(self, conn: sqlite3.Connection) -> None:
|
||||
"""A valid select SQL creates the named view."""
|
||||
_create_view_safe(conn, "test_view", "SELECT 1 AS val")
|
||||
views = _get_names(conn, "view")
|
||||
assert "test_view" in views
|
||||
|
||||
def test_replaces_existing_view(self, conn: sqlite3.Connection) -> None:
|
||||
"""Calling again with a new SQL replaces the existing view."""
|
||||
_create_view_safe(conn, "test_view", "SELECT 1 AS val")
|
||||
_create_view_safe(conn, "test_view", "SELECT 2 AS val")
|
||||
result = conn.execute("SELECT val FROM test_view").fetchone()
|
||||
assert result == (2,)
|
||||
|
||||
def test_logs_warning_on_sqlite_error(
|
||||
self,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""sqlite3.Error during CREATE VIEW logs a warning instead of raising."""
|
||||
mock_conn = MagicMock()
|
||||
mock_conn.execute.side_effect = [
|
||||
None,
|
||||
sqlite3.OperationalError("parse error"),
|
||||
]
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
_create_view_safe(mock_conn, "bad_view", "SELECT 1")
|
||||
assert "Skipping view bad_view" in caplog.text
|
||||
assert "parse error" in caplog.text
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestGrafanaViews
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestGrafanaViews:
|
||||
"""Tests for create_grafana_views and individual view builders."""
|
||||
|
||||
def test_all_views_created_with_full_schema(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""All 13 Grafana views are created when all source tables are present."""
|
||||
_make_rates_table(conn)
|
||||
_make_ticks_table(conn)
|
||||
_make_history_deals_full(conn)
|
||||
_make_history_orders_table(conn)
|
||||
create_snapshot_tables(conn)
|
||||
create_grafana_views(conn)
|
||||
views = _get_names(conn, "view")
|
||||
expected = {
|
||||
"grafana_rates",
|
||||
"grafana_ticks",
|
||||
"grafana_history_deals",
|
||||
"grafana_history_orders",
|
||||
"grafana_trade_deals",
|
||||
"grafana_cash_events",
|
||||
"grafana_realized_pnl",
|
||||
"grafana_symbol_pnl",
|
||||
"grafana_trade_stats",
|
||||
"grafana_account_snapshots",
|
||||
"grafana_position_snapshots",
|
||||
"grafana_order_snapshots",
|
||||
"grafana_terminal_snapshots",
|
||||
}
|
||||
assert expected.issubset(views)
|
||||
|
||||
def test_stale_view_dropped_when_source_table_disappears(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""create_grafana_views drops a previously created view whose source is gone."""
|
||||
_make_ticks_table(conn)
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_ticks" in _get_names(conn, "view")
|
||||
conn.execute("DROP TABLE ticks")
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_ticks" not in _get_names(conn, "view")
|
||||
|
||||
def test_grafana_rates_skipped_when_table_absent(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_rates is skipped when rates table is missing."""
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_rates" not in _get_names(conn, "view")
|
||||
|
||||
def test_grafana_rates_skipped_when_required_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_rates is skipped when rates table lacks required columns."""
|
||||
conn.execute("CREATE TABLE rates (open REAL)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_rates" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_rates" in caplog.text
|
||||
|
||||
def test_grafana_ticks_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_ticks is skipped when ticks table lacks required columns."""
|
||||
conn.execute("CREATE TABLE ticks (bid REAL)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_ticks" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_ticks" in caplog.text
|
||||
|
||||
def test_grafana_history_deals_skipped_when_time_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_history_deals is skipped when history_deals.time is missing."""
|
||||
conn.execute("CREATE TABLE history_deals (symbol TEXT)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_history_deals" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_history_deals" in caplog.text
|
||||
|
||||
def test_grafana_history_orders_skipped_when_time_setup_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_history_orders is skipped when time_setup is absent."""
|
||||
conn.execute("CREATE TABLE history_orders (symbol TEXT)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_history_orders" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_history_orders" in caplog.text
|
||||
|
||||
def test_grafana_trade_deals_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_trade_deals is skipped when history_deals missing time/type."""
|
||||
conn.execute("CREATE TABLE history_deals (symbol TEXT)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_trade_deals" not in _get_names(conn, "view")
|
||||
|
||||
def test_grafana_cash_events_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_cash_events is skipped when history_deals missing time/type."""
|
||||
conn.execute("CREATE TABLE history_deals (symbol TEXT)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_cash_events" not in _get_names(conn, "view")
|
||||
|
||||
def test_grafana_realized_pnl_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_realized_pnl is skipped when history_deals missing required cols."""
|
||||
conn.execute("CREATE TABLE history_deals (time TEXT, type INTEGER)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_realized_pnl" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_realized_pnl" in caplog.text
|
||||
|
||||
def test_grafana_realized_pnl_skipped_when_entry_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_realized_pnl is skipped when entry column is absent."""
|
||||
_make_history_deals_minimal(conn)
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_realized_pnl" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_realized_pnl" in caplog.text
|
||||
|
||||
def test_grafana_symbol_pnl_skipped_when_required_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_symbol_pnl is skipped when required columns are absent."""
|
||||
conn.execute("CREATE TABLE history_deals (time TEXT, type INTEGER)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_symbol_pnl" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_symbol_pnl" in caplog.text
|
||||
|
||||
def test_grafana_symbol_pnl_without_volume_and_price(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""grafana_symbol_pnl is created with only required columns."""
|
||||
conn.execute(
|
||||
"CREATE TABLE history_deals"
|
||||
" (time TEXT, symbol TEXT, profit REAL, type INTEGER, entry INTEGER)"
|
||||
)
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_symbol_pnl" in _get_names(conn, "view")
|
||||
|
||||
def test_grafana_symbol_pnl_with_volume_and_price(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""grafana_symbol_pnl includes volume and price columns when present."""
|
||||
_make_history_deals_full(conn)
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_symbol_pnl" in _get_names(conn, "view")
|
||||
# View columns include volume and price
|
||||
cols = {row[1] for row in conn.execute("PRAGMA table_info(grafana_symbol_pnl)")}
|
||||
assert "volume" in cols
|
||||
assert "price" in cols
|
||||
|
||||
def test_grafana_trade_stats_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""grafana_trade_stats is skipped when history_deals missing required cols."""
|
||||
conn.execute("CREATE TABLE history_deals (time TEXT)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_trade_stats" not in _get_names(conn, "view")
|
||||
assert "Skipping grafana_trade_stats" in caplog.text
|
||||
|
||||
def test_grafana_trade_stats_without_entry_col(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""grafana_trade_stats is a static summary view with no time column."""
|
||||
_make_history_deals_minimal(conn)
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_trade_stats" in _get_names(conn, "view")
|
||||
cols = {
|
||||
row[1] for row in conn.execute("PRAGMA table_info(grafana_trade_stats)")
|
||||
}
|
||||
assert "time" not in cols
|
||||
assert "symbol" in cols
|
||||
|
||||
def test_grafana_trade_stats_with_entry_col(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""grafana_trade_stats is a static summary view with no time column."""
|
||||
_make_history_deals_full(conn)
|
||||
create_grafana_views(conn)
|
||||
assert "grafana_trade_stats" in _get_names(conn, "view")
|
||||
cols = {
|
||||
row[1] for row in conn.execute("PRAGMA table_info(grafana_trade_stats)")
|
||||
}
|
||||
assert "time" not in cols
|
||||
assert "symbol" in cols
|
||||
|
||||
def test_snapshot_views_skipped_when_snapshot_tables_absent(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""Snapshot views are skipped when snapshot tables are not created."""
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
create_grafana_views(conn)
|
||||
views = _get_names(conn, "view")
|
||||
assert "grafana_account_snapshots" not in views
|
||||
assert "grafana_position_snapshots" not in views
|
||||
assert "grafana_order_snapshots" not in views
|
||||
assert "grafana_terminal_snapshots" not in views
|
||||
|
||||
def test_build_snapshot_view_with_only_run_id_col(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""_build_snapshot_view exposes time and run_id when table has only run_id."""
|
||||
create_snapshot_tables(conn)
|
||||
conn.execute("CREATE TABLE only_run (run_id INTEGER NOT NULL)")
|
||||
run_id = start_snapshot_run(conn, 1000)
|
||||
record_snapshot_run(conn, run_id, "ok")
|
||||
conn.execute("INSERT INTO only_run (run_id) VALUES (?)", (run_id,))
|
||||
_build_snapshot_view(conn, "test_view", "only_run")
|
||||
assert "test_view" in _get_names(conn, "view")
|
||||
cols = {row[1] for row in conn.execute("PRAGMA table_info(test_view)")}
|
||||
assert "time" in cols
|
||||
assert "run_id" in cols
|
||||
|
||||
def test_build_snapshot_view_skips_when_snapshot_runs_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""_build_snapshot_view skips view when snapshot_runs has wrong columns."""
|
||||
conn.execute("CREATE TABLE only_run (run_id INTEGER NOT NULL)")
|
||||
conn.execute("CREATE TABLE snapshot_runs (foo TEXT)")
|
||||
_build_snapshot_view(conn, "test_view", "only_run")
|
||||
views = _get_names(conn, "view")
|
||||
assert "test_view" not in views
|
||||
|
||||
def test_build_snapshot_view_skips_when_run_id_col_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""_build_snapshot_view skips view when the table lacks run_id."""
|
||||
create_snapshot_tables(conn)
|
||||
conn.execute("CREATE TABLE no_run_id (symbol TEXT)")
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.grafana"):
|
||||
_build_snapshot_view(conn, "test_view", "no_run_id")
|
||||
assert "test_view" not in _get_names(conn, "view")
|
||||
assert "missing run_id column" in caplog.text
|
||||
|
||||
def test_snapshot_view_excludes_failed_run_rows(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Snapshot views hide rows from failed runs."""
|
||||
create_snapshot_tables(conn)
|
||||
run_id = start_snapshot_run(conn, 1000)
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots"
|
||||
" (run_id, login, balance, equity, margin, margin_free, profit)"
|
||||
" VALUES (?, 12345, 10000.0, 9800.0, 200.0, 9600.0, -200.0)",
|
||||
(run_id,),
|
||||
)
|
||||
record_snapshot_run(conn, run_id, "error", "terminal offline")
|
||||
create_grafana_views(conn)
|
||||
rows = conn.execute("SELECT * FROM grafana_account_snapshots").fetchall()
|
||||
assert rows == []
|
||||
|
||||
def test_snapshot_view_includes_ok_run_rows(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Snapshot views show rows from successful runs and expose run_id."""
|
||||
create_snapshot_tables(conn)
|
||||
run_id = start_snapshot_run(conn, 2000)
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots"
|
||||
" (run_id, login, balance, equity, margin, margin_free, profit)"
|
||||
" VALUES (?, 12345, 10000.0, 9800.0, 200.0, 9600.0, -200.0)",
|
||||
(run_id,),
|
||||
)
|
||||
record_snapshot_run(conn, run_id, "ok")
|
||||
create_grafana_views(conn)
|
||||
rows = conn.execute(
|
||||
"SELECT time, run_id, login FROM grafana_account_snapshots"
|
||||
).fetchall()
|
||||
assert rows == [(2000, run_id, 12345)]
|
||||
cols = {
|
||||
row[1]
|
||||
for row in conn.execute("PRAGMA table_info(grafana_account_snapshots)")
|
||||
}
|
||||
assert "run_id" in cols
|
||||
|
||||
def test_snapshot_view_same_second_ok_and_error_no_cross_contamination(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""An ok and error run sharing observed_at expose only the ok run's rows."""
|
||||
create_snapshot_tables(conn)
|
||||
run_err = start_snapshot_run(conn, 3000)
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots (run_id, login) VALUES (?, 99)",
|
||||
(run_err,),
|
||||
)
|
||||
record_snapshot_run(conn, run_err, "error")
|
||||
run_ok = start_snapshot_run(conn, 3000)
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots (run_id, login) VALUES (?, 12345)",
|
||||
(run_ok,),
|
||||
)
|
||||
record_snapshot_run(conn, run_ok, "ok")
|
||||
create_grafana_views(conn)
|
||||
rows = conn.execute("SELECT login FROM grafana_account_snapshots").fetchall()
|
||||
assert rows == [(12345,)]
|
||||
|
||||
def test_snapshot_view_two_ok_runs_same_second_no_duplication(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Two ok runs sharing observed_at each produce exactly one row in the view."""
|
||||
create_snapshot_tables(conn)
|
||||
run1 = start_snapshot_run(conn, 4000)
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots (run_id, login) VALUES (?, 1)",
|
||||
(run1,),
|
||||
)
|
||||
record_snapshot_run(conn, run1, "ok")
|
||||
run2 = start_snapshot_run(conn, 4000)
|
||||
conn.execute(
|
||||
"INSERT INTO account_snapshots (run_id, login) VALUES (?, 2)",
|
||||
(run2,),
|
||||
)
|
||||
record_snapshot_run(conn, run2, "ok")
|
||||
create_grafana_views(conn)
|
||||
rows = conn.execute("SELECT login FROM grafana_account_snapshots").fetchall()
|
||||
assert len(rows) == 2
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestGrafanaIndexes
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestGrafanaIndexes:
|
||||
"""Tests for create_grafana_indexes."""
|
||||
|
||||
def test_all_indexes_created_with_full_schema(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""All 9 indexes are created when all source tables are present."""
|
||||
_make_rates_table(conn)
|
||||
_make_ticks_table(conn)
|
||||
_make_history_deals_full(conn)
|
||||
_make_history_orders_table(conn)
|
||||
create_snapshot_tables(conn)
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_rates_time_symbol_timeframe" in indexes
|
||||
assert "idx_ticks_time_symbol" in indexes
|
||||
assert "idx_history_deals_time_symbol" in indexes
|
||||
assert "idx_history_deals_symbol_time" in indexes
|
||||
assert "idx_history_orders_time_setup_symbol" in indexes
|
||||
assert "idx_account_snapshots_time_login" in indexes
|
||||
assert "idx_position_snapshots_time_symbol" in indexes
|
||||
assert "idx_order_snapshots_time_symbol" in indexes
|
||||
assert "idx_snapshot_runs_time_status" in indexes
|
||||
|
||||
def test_no_indexes_created_when_tables_absent(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""No indexes are created when tables are absent."""
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert not any(name.startswith("idx_") for name in indexes)
|
||||
|
||||
def test_indexes_for_snapshot_tables_skipped_when_absent(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Snapshot table indexes are skipped when snapshot tables don't exist."""
|
||||
_make_history_deals_full(conn)
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_account_snapshots_time_login" not in indexes
|
||||
assert "idx_position_snapshots_time_symbol" not in indexes
|
||||
assert "idx_order_snapshots_time_symbol" not in indexes
|
||||
assert "idx_snapshot_runs_time_status" not in indexes
|
||||
|
||||
def test_rates_index_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Rates index is skipped when required columns are absent."""
|
||||
conn.execute("CREATE TABLE rates (open REAL)")
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_rates_time_symbol_timeframe" not in indexes
|
||||
|
||||
def test_ticks_index_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Ticks index is skipped when required columns are absent."""
|
||||
conn.execute("CREATE TABLE ticks (bid REAL)")
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_ticks_time_symbol" not in indexes
|
||||
|
||||
def test_deals_indexes_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""history_deals indexes are skipped when required columns are absent."""
|
||||
conn.execute("CREATE TABLE history_deals (ticket INTEGER)")
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_history_deals_time_symbol" not in indexes
|
||||
|
||||
def test_orders_index_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""history_orders index is skipped when required columns are absent."""
|
||||
conn.execute("CREATE TABLE history_orders (ticket INTEGER)")
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_history_orders_time_setup_symbol" not in indexes
|
||||
|
||||
def test_snapshot_indexes_skipped_when_cols_missing(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""Snapshot table indexes are skipped when required columns are absent."""
|
||||
conn.execute("CREATE TABLE account_snapshots (foo TEXT)")
|
||||
conn.execute("CREATE TABLE position_snapshots (foo TEXT)")
|
||||
conn.execute("CREATE TABLE order_snapshots (foo TEXT)")
|
||||
conn.execute("CREATE TABLE snapshot_runs (foo TEXT)")
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_account_snapshots_time_login" not in indexes
|
||||
assert "idx_position_snapshots_time_symbol" not in indexes
|
||||
assert "idx_order_snapshots_time_symbol" not in indexes
|
||||
assert "idx_snapshot_runs_time_status" not in indexes
|
||||
|
||||
def test_indexes_are_idempotent(self, conn: sqlite3.Connection) -> None:
|
||||
"""Creating indexes twice does not raise (IF NOT EXISTS)."""
|
||||
_make_rates_table(conn)
|
||||
create_grafana_indexes(conn)
|
||||
create_grafana_indexes(conn)
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_rates_time_symbol_timeframe" in indexes
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestEnsureGrafanaSchema
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestEnsureGrafanaSchema:
|
||||
"""Tests for ensure_grafana_schema."""
|
||||
|
||||
def test_creates_all_tables_views_and_indexes(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""ensure_grafana_schema creates snapshot tables, views, and indexes."""
|
||||
_make_rates_table(conn)
|
||||
_make_history_deals_full(conn)
|
||||
ensure_grafana_schema(conn)
|
||||
tables = _get_names(conn, "table")
|
||||
assert "snapshot_runs" in tables
|
||||
assert "account_snapshots" in tables
|
||||
views = _get_names(conn, "view")
|
||||
assert "grafana_rates" in views
|
||||
assert "grafana_account_snapshots" in views
|
||||
indexes = _get_names(conn, "index")
|
||||
assert "idx_rates_time_symbol_timeframe" in indexes
|
||||
|
||||
def test_is_idempotent(self, conn: sqlite3.Connection) -> None:
|
||||
"""Calling ensure_grafana_schema twice does not raise."""
|
||||
ensure_grafana_schema(conn)
|
||||
ensure_grafana_schema(conn)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestSnapshotInserts
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestSnapshotInserts:
|
||||
"""Tests for snapshot insert helpers."""
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def setup_tables(self, conn: sqlite3.Connection) -> None:
|
||||
"""Create snapshot tables before each insert test."""
|
||||
create_snapshot_tables(conn)
|
||||
|
||||
def test_insert_account_snapshot(self, conn: sqlite3.Connection) -> None:
|
||||
"""insert_account_snapshot appends a row with correct values."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
row: dict[str, object] = {
|
||||
"login": 12345,
|
||||
"currency": "USD",
|
||||
"balance": 10000.0,
|
||||
"equity": 9800.0,
|
||||
"margin": 200.0,
|
||||
"margin_free": 9800.0,
|
||||
"margin_level": 4900.0,
|
||||
"profit": -200.0,
|
||||
"leverage": 100,
|
||||
}
|
||||
insert_account_snapshot(conn, run_id, row)
|
||||
result = conn.execute(
|
||||
"SELECT login, currency, balance FROM account_snapshots"
|
||||
).fetchone()
|
||||
assert result == (12345, "USD", 10000.0)
|
||||
|
||||
def test_insert_account_snapshot_partial_row(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_account_snapshot works when some fields are missing (uses None)."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
insert_account_snapshot(conn, run_id, {"login": 1})
|
||||
result = conn.execute(
|
||||
"SELECT login, currency FROM account_snapshots"
|
||||
).fetchone()
|
||||
assert result == (1, None)
|
||||
|
||||
def test_insert_position_snapshots_with_rows(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_position_snapshots appends each position row."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
rows: list[dict[str, object]] = [
|
||||
{"ticket": 1, "symbol": "EURUSD", "volume": 0.1, "profit": 10.0},
|
||||
{"ticket": 2, "symbol": "GBPUSD", "volume": 0.2, "profit": -5.0},
|
||||
]
|
||||
insert_position_snapshots(conn, run_id, 12345, rows)
|
||||
count = conn.execute("SELECT COUNT(*) FROM position_snapshots").fetchone()[0]
|
||||
assert count == 2
|
||||
|
||||
def test_insert_position_snapshots_noop_when_empty(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_position_snapshots is a no-op when rows is empty."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
insert_position_snapshots(conn, run_id, 12345, [])
|
||||
count = conn.execute("SELECT COUNT(*) FROM position_snapshots").fetchone()[0]
|
||||
assert count == 0
|
||||
|
||||
def test_insert_order_snapshots_with_rows(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_order_snapshots appends each order row."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
rows: list[dict[str, object]] = [
|
||||
{"ticket": 10, "symbol": "EURUSD", "type": 2, "volume_current": 0.1},
|
||||
]
|
||||
insert_order_snapshots(conn, run_id, 12345, rows)
|
||||
count = conn.execute("SELECT COUNT(*) FROM order_snapshots").fetchone()[0]
|
||||
assert count == 1
|
||||
|
||||
def test_insert_order_snapshots_noop_when_empty(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_order_snapshots is a no-op when rows is empty."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
insert_order_snapshots(conn, run_id, 12345, [])
|
||||
count = conn.execute("SELECT COUNT(*) FROM order_snapshots").fetchone()[0]
|
||||
assert count == 0
|
||||
|
||||
def test_insert_order_snapshots_normalizes_timestamp_time_setup(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_order_snapshots converts pd.Timestamp time_setup to epoch int."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
ts = pd.Timestamp("2024-01-15 10:30:00", tz="UTC")
|
||||
rows: list[dict[str, object]] = [{"ticket": 10, "time_setup": ts}]
|
||||
insert_order_snapshots(conn, run_id, 12345, rows)
|
||||
stored = conn.execute("SELECT time_setup FROM order_snapshots").fetchone()[0]
|
||||
assert stored == int(ts.timestamp())
|
||||
|
||||
def test_insert_order_snapshots_stores_int_time_setup(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_order_snapshots stores an integer time_setup as-is."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
rows: list[dict[str, object]] = [{"ticket": 10, "time_setup": 1705314600}]
|
||||
insert_order_snapshots(conn, run_id, 12345, rows)
|
||||
stored = conn.execute("SELECT time_setup FROM order_snapshots").fetchone()[0]
|
||||
assert stored == 1705314600
|
||||
|
||||
def test_insert_order_snapshots_stores_null_for_unknown_time_setup_type(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""insert_order_snapshots stores NULL for an unrecognized time_setup type."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
rows: list[dict[str, object]] = [{"ticket": 10, "time_setup": "not_a_time"}]
|
||||
insert_order_snapshots(conn, run_id, 12345, rows)
|
||||
stored = conn.execute("SELECT time_setup FROM order_snapshots").fetchone()[0]
|
||||
assert stored is None
|
||||
|
||||
def test_insert_terminal_snapshot(self, conn: sqlite3.Connection) -> None:
|
||||
"""insert_terminal_snapshot appends a terminal info row."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
row: dict[str, object] = {
|
||||
"name": "MetaTrader 5",
|
||||
"connected": 1,
|
||||
"community_account": 0,
|
||||
"trade_allowed": 1,
|
||||
"trade_expert": 1,
|
||||
"path": "/mt5",
|
||||
"company": "Broker",
|
||||
"language": "en",
|
||||
}
|
||||
insert_terminal_snapshot(conn, run_id, row)
|
||||
result = conn.execute(
|
||||
"SELECT name, connected FROM terminal_snapshots"
|
||||
).fetchone()
|
||||
assert result == ("MetaTrader 5", 1)
|
||||
|
||||
def test_start_snapshot_run_returns_incrementing_ids(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""start_snapshot_run returns a unique run_id for each call."""
|
||||
run1 = start_snapshot_run(conn, 1700000000)
|
||||
run2 = start_snapshot_run(conn, 1700000000)
|
||||
assert run1 != run2
|
||||
|
||||
def test_record_snapshot_run_with_detail(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""record_snapshot_run stores status and detail text."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
record_snapshot_run(conn, run_id, "error", "RuntimeError: boom")
|
||||
row = conn.execute("SELECT status, detail FROM snapshot_runs").fetchone()
|
||||
assert row == ("error", "RuntimeError: boom")
|
||||
|
||||
def test_record_snapshot_run_without_detail(
|
||||
self,
|
||||
conn: sqlite3.Connection,
|
||||
) -> None:
|
||||
"""record_snapshot_run stores None for detail when omitted."""
|
||||
run_id = start_snapshot_run(conn, 1700000000)
|
||||
record_snapshot_run(conn, run_id, "ok")
|
||||
row = conn.execute("SELECT status, detail FROM snapshot_runs").fetchone()
|
||||
assert row == ("ok", None)
|
||||
@@ -61,6 +61,8 @@ from mt5cli.sdk import (
|
||||
terminal_info,
|
||||
update_history,
|
||||
update_history_with_config,
|
||||
update_observability,
|
||||
update_observability_with_config,
|
||||
version,
|
||||
)
|
||||
from mt5cli.utils import Dataset, IfExists, coerce_login
|
||||
@@ -2746,3 +2748,375 @@ class TestSubstituteMappingValues:
|
||||
result = substitute_mapping_values(data, keys={"mt5_login"})
|
||||
# tuple is returned as-is; inner dict is NOT visited
|
||||
assert result == {"accounts": ({"mt5_login": "${MT5_LOGIN}"},)}
|
||||
|
||||
|
||||
class TestUpdateObservability:
|
||||
"""Tests for update_observability and update_observability_with_config."""
|
||||
|
||||
@pytest.fixture
|
||||
def mock_client(self) -> MagicMock:
|
||||
"""Mock client returning minimal valid frames."""
|
||||
client = MagicMock()
|
||||
client.account_info_as_df.return_value = pd.DataFrame([
|
||||
{
|
||||
"login": 12345,
|
||||
"currency": "USD",
|
||||
"balance": 10000.0,
|
||||
"equity": 10000.0,
|
||||
"margin": 0.0,
|
||||
"margin_free": 10000.0,
|
||||
"margin_level": 0.0,
|
||||
"profit": 0.0,
|
||||
"leverage": 100,
|
||||
}
|
||||
])
|
||||
client.positions_get_as_df.return_value = pd.DataFrame()
|
||||
client.orders_get_as_df.return_value = pd.DataFrame()
|
||||
client.terminal_info_as_df.return_value = pd.DataFrame([
|
||||
{
|
||||
"name": "MetaTrader 5",
|
||||
"connected": 1,
|
||||
"community_account": 0,
|
||||
"trade_allowed": 1,
|
||||
"trade_expert": 1,
|
||||
"path": "/mt5",
|
||||
"company": "Broker",
|
||||
"language": "en",
|
||||
}
|
||||
])
|
||||
return client
|
||||
|
||||
def test_update_observability_creates_snapshot_tables(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Snapshot tables are created in the output database."""
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=mock_client, output=output)
|
||||
with sqlite3.connect(output) as conn:
|
||||
tables = {
|
||||
row[0]
|
||||
for row in conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type='table'"
|
||||
).fetchall()
|
||||
}
|
||||
assert "snapshot_runs" in tables
|
||||
assert "account_snapshots" in tables
|
||||
assert "position_snapshots" in tables
|
||||
|
||||
def test_update_observability_records_ok_on_success(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""snapshot_runs records 'ok' status on a successful run."""
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=mock_client, output=output)
|
||||
with sqlite3.connect(output) as conn:
|
||||
row = conn.execute("SELECT status FROM snapshot_runs").fetchone()
|
||||
assert row == ("ok",)
|
||||
|
||||
def test_update_observability_records_error_on_failure(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""snapshot_runs records 'error' and re-raises when a snapshot fails."""
|
||||
mock_client.account_info_as_df.side_effect = RuntimeError("boom")
|
||||
output = tmp_path / "obs.db"
|
||||
with pytest.raises(RuntimeError, match="boom"):
|
||||
update_observability(client=mock_client, output=output)
|
||||
with sqlite3.connect(output) as conn:
|
||||
row = conn.execute("SELECT status FROM snapshot_runs").fetchone()
|
||||
assert row == ("error",)
|
||||
|
||||
def test_update_observability_skips_ensure_grafana_schema_when_disabled(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""with_grafana_schema=False does not call ensure_grafana_schema."""
|
||||
spy = mocker.spy(sdk, "ensure_grafana_schema")
|
||||
update_observability(
|
||||
client=mock_client,
|
||||
output=tmp_path / "obs.db",
|
||||
with_grafana_schema=False,
|
||||
)
|
||||
spy.assert_not_called()
|
||||
|
||||
def test_update_observability_calls_ensure_grafana_schema_by_default(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""with_grafana_schema=True calls ensure_grafana_schema."""
|
||||
spy = mocker.spy(sdk, "ensure_grafana_schema")
|
||||
update_observability(
|
||||
client=mock_client,
|
||||
output=tmp_path / "obs.db",
|
||||
with_grafana_schema=True,
|
||||
)
|
||||
spy.assert_called_once()
|
||||
|
||||
def test_update_observability_skips_account_when_disabled(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""include_account=False does not call account_info_as_df."""
|
||||
update_observability(
|
||||
client=mock_client,
|
||||
output=tmp_path / "obs.db",
|
||||
include_account=False,
|
||||
)
|
||||
mock_client.account_info_as_df.assert_not_called()
|
||||
|
||||
def test_update_observability_skips_positions_when_disabled(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""include_positions=False does not call positions_get_as_df."""
|
||||
update_observability(
|
||||
client=mock_client,
|
||||
output=tmp_path / "obs.db",
|
||||
include_positions=False,
|
||||
)
|
||||
mock_client.positions_get_as_df.assert_not_called()
|
||||
|
||||
def test_update_observability_skips_orders_when_disabled(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""include_orders=False does not call orders_get_as_df."""
|
||||
update_observability(
|
||||
client=mock_client,
|
||||
output=tmp_path / "obs.db",
|
||||
include_orders=False,
|
||||
)
|
||||
mock_client.orders_get_as_df.assert_not_called()
|
||||
|
||||
def test_update_observability_skips_terminal_when_disabled(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""include_terminal=False does not call terminal_info_as_df."""
|
||||
update_observability(
|
||||
client=mock_client,
|
||||
output=tmp_path / "obs.db",
|
||||
include_terminal=False,
|
||||
)
|
||||
mock_client.terminal_info_as_df.assert_not_called()
|
||||
|
||||
def test_update_observability_with_positions_rows(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Non-empty positions are written to position_snapshots."""
|
||||
mock_client.positions_get_as_df.return_value = pd.DataFrame([
|
||||
{
|
||||
"ticket": 1,
|
||||
"position_id": 1,
|
||||
"symbol": "EURUSD",
|
||||
"type": 0,
|
||||
"volume": 0.1,
|
||||
"price_open": 1.1,
|
||||
"price_current": 1.1,
|
||||
"profit": 0.0,
|
||||
"swap": 0.0,
|
||||
"comment": "",
|
||||
"magic": 0,
|
||||
}
|
||||
])
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=mock_client, output=output)
|
||||
with sqlite3.connect(output) as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM position_snapshots").fetchone()[
|
||||
0
|
||||
]
|
||||
assert count == 1
|
||||
|
||||
def test_update_observability_with_order_rows(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Non-empty orders are written to order_snapshots."""
|
||||
mock_client.orders_get_as_df.return_value = pd.DataFrame([
|
||||
{
|
||||
"ticket": 10,
|
||||
"symbol": "EURUSD",
|
||||
"type": 2,
|
||||
"volume_current": 0.1,
|
||||
"price_open": 1.2,
|
||||
"price_current": 1.1,
|
||||
"state": 1,
|
||||
"comment": "",
|
||||
"magic": 0,
|
||||
"time_setup": 1700000000,
|
||||
}
|
||||
])
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=mock_client, output=output)
|
||||
with sqlite3.connect(output) as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM order_snapshots").fetchone()[0]
|
||||
assert count == 1
|
||||
|
||||
def test_update_observability_symbol_filter_positions(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Symbol filter fetches all positions in one call and filters client-side."""
|
||||
client = MagicMock()
|
||||
client.account_info_as_df.return_value = pd.DataFrame([{"login": 1}])
|
||||
# All positions; only EURUSD matches the filter
|
||||
client.positions_get_as_df.return_value = pd.DataFrame([
|
||||
{"ticket": 1, "symbol": "EURUSD", "volume": 0.1, "profit": 0.0},
|
||||
{"ticket": 2, "symbol": "USDJPY", "volume": 0.2, "profit": 0.0},
|
||||
])
|
||||
client.orders_get_as_df.return_value = pd.DataFrame()
|
||||
client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=client, output=output, symbols=["EURUSD", "GBPUSD"])
|
||||
assert client.positions_get_as_df.call_count == 1
|
||||
with sqlite3.connect(output) as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM position_snapshots").fetchone()[
|
||||
0
|
||||
]
|
||||
assert count == 1
|
||||
|
||||
def test_update_observability_symbol_filter_orders(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Symbol filter fetches all orders in one call and filters client-side."""
|
||||
client = MagicMock()
|
||||
client.account_info_as_df.return_value = pd.DataFrame([{"login": 1}])
|
||||
client.positions_get_as_df.return_value = pd.DataFrame()
|
||||
# All orders; only EURUSD matches the filter
|
||||
client.orders_get_as_df.return_value = pd.DataFrame([
|
||||
{"ticket": 10, "symbol": "EURUSD", "volume_current": 0.1},
|
||||
{"ticket": 11, "symbol": "USDJPY", "volume_current": 0.5},
|
||||
])
|
||||
client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=client, output=output, symbols=["EURUSD", "GBPUSD"])
|
||||
assert client.orders_get_as_df.call_count == 1
|
||||
with sqlite3.connect(output) as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM order_snapshots").fetchone()[0]
|
||||
assert count == 1
|
||||
|
||||
def test_update_observability_symbol_filter_no_symbol_col(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Symbol filter is skipped when positions df has no symbol column."""
|
||||
client = MagicMock()
|
||||
client.account_info_as_df.return_value = pd.DataFrame([{"login": 1}])
|
||||
# No symbol column in positions — all rows pass through unfiltered
|
||||
client.positions_get_as_df.return_value = pd.DataFrame([
|
||||
{"ticket": 1, "volume": 0.1},
|
||||
])
|
||||
client.orders_get_as_df.return_value = pd.DataFrame()
|
||||
client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=client, output=output, symbols=["EURUSD"])
|
||||
with sqlite3.connect(output) as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM position_snapshots").fetchone()[
|
||||
0
|
||||
]
|
||||
assert count == 1
|
||||
|
||||
def test_update_observability_account_none_login(
|
||||
self,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Account row with no login key returns None login for downstream helpers."""
|
||||
client = MagicMock()
|
||||
client.account_info_as_df.return_value = pd.DataFrame([{"balance": 10000.0}])
|
||||
client.positions_get_as_df.return_value = pd.DataFrame()
|
||||
client.orders_get_as_df.return_value = pd.DataFrame()
|
||||
client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
output = tmp_path / "obs.db"
|
||||
update_observability(client=client, output=output)
|
||||
with sqlite3.connect(output) as conn:
|
||||
row = conn.execute("SELECT login FROM account_snapshots").fetchone()
|
||||
assert row is not None
|
||||
assert row[0] is None
|
||||
|
||||
def test_update_observability_empty_account_logs_warning(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""Empty account_info_as_df logs a warning and does not write account row."""
|
||||
mock_client.account_info_as_df.return_value = pd.DataFrame()
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.sdk"):
|
||||
update_observability(client=mock_client, output=tmp_path / "obs.db")
|
||||
assert "account_info_as_df returned empty frame" in caplog.text
|
||||
with sqlite3.connect(tmp_path / "obs.db") as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM account_snapshots").fetchone()[0]
|
||||
assert count == 0
|
||||
|
||||
def test_update_observability_empty_terminal_logs_warning(
|
||||
self,
|
||||
mock_client: MagicMock,
|
||||
tmp_path: Path,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""Empty terminal_info_as_df logs a warning and does not write terminal row."""
|
||||
mock_client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
with caplog.at_level(logging.WARNING, logger="mt5cli.sdk"):
|
||||
update_observability(client=mock_client, output=tmp_path / "obs.db")
|
||||
assert "terminal_info_as_df returned empty frame" in caplog.text
|
||||
with sqlite3.connect(tmp_path / "obs.db") as conn:
|
||||
count = conn.execute("SELECT COUNT(*) FROM terminal_snapshots").fetchone()[
|
||||
0
|
||||
]
|
||||
assert count == 0
|
||||
|
||||
def test_update_observability_with_config_opens_and_closes_connection(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""update_observability_with_config manages the MT5 connection lifecycle."""
|
||||
mock_client = MagicMock()
|
||||
mock_client.account_info_as_df.return_value = pd.DataFrame()
|
||||
mock_client.positions_get_as_df.return_value = pd.DataFrame()
|
||||
mock_client.orders_get_as_df.return_value = pd.DataFrame()
|
||||
mock_client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=mock_client)
|
||||
update_observability_with_config(output=tmp_path / "obs.db")
|
||||
mock_client.initialize_and_login_mt5.assert_called_once()
|
||||
mock_client.shutdown.assert_called_once()
|
||||
|
||||
def test_update_observability_with_config_passes_symbols(
|
||||
self,
|
||||
mocker: MockerFixture,
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""update_observability_with_config forwards symbols to update_observability."""
|
||||
mock_client = MagicMock()
|
||||
mock_client.account_info_as_df.return_value = pd.DataFrame()
|
||||
mock_client.positions_get_as_df.return_value = pd.DataFrame()
|
||||
mock_client.orders_get_as_df.return_value = pd.DataFrame()
|
||||
mock_client.terminal_info_as_df.return_value = pd.DataFrame()
|
||||
mocker.patch("mt5cli.sdk.Mt5DataClient", return_value=mock_client)
|
||||
spy = mocker.patch("mt5cli.sdk.update_observability")
|
||||
update_observability_with_config(
|
||||
output=tmp_path / "obs.db",
|
||||
symbols=["EURUSD"],
|
||||
include_account=False,
|
||||
)
|
||||
spy.assert_called_once()
|
||||
call_kwargs = spy.call_args.kwargs
|
||||
assert call_kwargs["symbols"] == ["EURUSD"]
|
||||
assert call_kwargs["include_account"] is False
|
||||
|
||||
Reference in New Issue
Block a user