From d27da02f3f305b5095c4ff53acee2a8b7ec9eff2 Mon Sep 17 00:00:00 2001 From: Daichi Narushima <1938249+dceoy@users.noreply.github.com> Date: Sun, 28 Jun 2026 07:55:26 +0900 Subject: [PATCH] feat: add Grafana-ready SQLite observability (#86) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * 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 * chore: apply mdformat to docs after grafana observability additions Co-Authored-By: Claude Sonnet 4.6 * chore: bump version to 1.1.0 Co-Authored-By: Claude Sonnet 4.6 * 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 * chore: apply ruff formatting and sync lock file Co-Authored-By: Claude Sonnet 4.6 * 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 * 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 * chore: format markdown tables Align table column widths in README and public-contract documentation. Co-Authored-By: Claude Haiku 4.5 * 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 * 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 * 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 * 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 * 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 --------- Co-authored-by: agent Co-authored-by: Claude Sonnet 4.6 --- README.md | 100 +++++ docs/api/public-contract.md | 83 +++- mt5cli/__init__.py | 4 + mt5cli/cli.py | 94 ++++ mt5cli/contract.py | 2 + mt5cli/grafana.py | 616 ++++++++++++++++++++++++++ mt5cli/sdk.py | 184 +++++++- mt5cli/utils.py | 3 +- pyproject.toml | 2 +- tests/conftest.py | 38 ++ tests/test_cli.py | 168 ++++++++ tests/test_grafana.py | 838 ++++++++++++++++++++++++++++++++++++ tests/test_sdk.py | 374 ++++++++++++++++ uv.lock | 2 +- 14 files changed, 2494 insertions(+), 14 deletions(-) create mode 100644 mt5cli/grafana.py create mode 100644 tests/test_grafana.py diff --git a/README.md b/README.md index 3cce189..5ac5317 100644 --- a/README.md +++ b/README.md @@ -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: diff --git a/docs/api/public-contract.md b/docs/api/public-contract.md index 134ad2a..9887e82 100644 --- a/docs/api/public-contract.md +++ b/docs/api/public-contract.md @@ -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 diff --git a/mt5cli/__init__.py b/mt5cli/__init__.py index ff552fa..ff6cb0e 100644 --- a/mt5cli/__init__.py +++ b/mt5cli/__init__.py @@ -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", ] diff --git a/mt5cli/cli.py b/mt5cli/cli.py index 7883876..88d5801 100644 --- a/mt5cli/cli.py +++ b/mt5cli/cli.py @@ -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() diff --git a/mt5cli/contract.py b/mt5cli/contract.py index 3d8567e..dcac983 100644 --- a/mt5cli/contract.py +++ b/mt5cli/contract.py @@ -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", }) diff --git a/mt5cli/grafana.py b/mt5cli/grafana.py new file mode 100644 index 0000000..b58c6a8 --- /dev/null +++ b/mt5cli/grafana.py @@ -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), + ) diff --git a/mt5cli/sdk.py b/mt5cli/sdk.py index d3646be..433d9ce 100644 --- a/mt5cli/sdk.py +++ b/mt5cli/sdk.py @@ -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, + ) diff --git a/mt5cli/utils.py b/mt5cli/utils.py index f43fe5e..8a779a2 100644 --- a/mt5cli/utils.py +++ b/mt5cli/utils.py @@ -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, diff --git a/pyproject.toml b/pyproject.toml index 9db47a0..9056a9e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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"}] diff --git a/tests/conftest.py b/tests/conftest.py index a30ef9b..85703f6 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -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) diff --git a/tests/test_cli.py b/tests/test_cli.py index 768539c..a60113b 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -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.""" diff --git a/tests/test_grafana.py b/tests/test_grafana.py new file mode 100644 index 0000000..fa2fb40 --- /dev/null +++ b/tests/test_grafana.py @@ -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) diff --git a/tests/test_sdk.py b/tests/test_sdk.py index 1401e1c..1b17fa9 100644 --- a/tests/test_sdk.py +++ b/tests/test_sdk.py @@ -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 diff --git a/uv.lock b/uv.lock index 2cb34a5..814cadc 100644 --- a/uv.lock +++ b/uv.lock @@ -487,7 +487,7 @@ wheels = [ [[package]] name = "mt5cli" -version = "1.0.3" +version = "1.1.0" source = { editable = "." } dependencies = [ { name = "click" },