Add incremental SQLite history SDK (#16)
* Add incremental SQLite history SDK for automated pipelines. Extract sqlite history helpers into a dedicated module and expose update_history APIs that resume from existing MAX(time) values instead of re-fetching fixed date ranges. Co-authored-by: Cursor <cursoragent@cursor.com> * Fix incremental history deals and stale rate view cleanup. Fetch account events once during incremental updates, drop stale rate_* views when timeframes change, and avoid SQLite variable limits on wide frames. Co-authored-by: Cursor <cursoragent@cursor.com> * Fix incremental deal filtering edge cases Co-authored-by: Cursor <cursoragent@cursor.com> * Address PR review feedback for incremental SQLite history. Make rate views collision-free, batch incremental resume queries, scope deduplication to appended boundaries, validate before opening MT5, use atomic SQLite transactions, and expand docs/tests for the new helpers. Co-authored-by: Cursor <cursoragent@cursor.com> * Document collect-history SQLite schema with ER diagram. Co-authored-by: Cursor <cursoragent@cursor.com> * Fix account-event filtering and drop legacy rates resume. Account events must follow only account_event_start, not per-symbol trade cursors. Require normalized rates schema and fail fast when timeframe is missing. Co-authored-by: Cursor <cursoragent@cursor.com> * Validate normalized rates schema before incremental resume. Require symbol, timeframe, and time on existing rates tables with clear ValueError messages, and add regression tests for malformed schemas. Co-authored-by: Cursor <cursoragent@cursor.com> --------- Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+222
-328
@@ -5,12 +5,23 @@ from __future__ import annotations
|
||||
import logging
|
||||
import sqlite3
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime
|
||||
from pathlib import Path # noqa: TC003
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import TYPE_CHECKING, Self, TypeVar
|
||||
|
||||
from pdmt5 import Mt5Config, Mt5DataClient
|
||||
|
||||
from .sqlite_history import (
|
||||
create_cash_events_view,
|
||||
create_history_indexes,
|
||||
create_positions_reconstructed_view,
|
||||
resolve_history_datasets,
|
||||
resolve_history_tick_flags,
|
||||
resolve_history_timeframes,
|
||||
write_collected_datasets,
|
||||
write_incremental_datasets,
|
||||
)
|
||||
from .utils import (
|
||||
Dataset,
|
||||
IfExists,
|
||||
@@ -20,7 +31,7 @@ from .utils import (
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import Callable, Iterator
|
||||
from collections.abc import Callable, Iterator, Sequence
|
||||
|
||||
import pandas as pd
|
||||
|
||||
@@ -48,22 +59,11 @@ __all__ = [
|
||||
"symbol_info_tick",
|
||||
"symbols",
|
||||
"terminal_info",
|
||||
"update_history",
|
||||
"update_history_with_config",
|
||||
"version",
|
||||
]
|
||||
|
||||
_TRADE_DEAL_TYPES: tuple[int, int] = (0, 1)
|
||||
_TRADE_DEAL_TYPES_SQL = f"({', '.join(str(value) for value in _TRADE_DEAL_TYPES)})"
|
||||
_POSITIONS_VIEW_REQUIRED_COLUMNS: frozenset[str] = frozenset({
|
||||
"position_id",
|
||||
"symbol",
|
||||
"time",
|
||||
"type",
|
||||
"entry",
|
||||
"volume",
|
||||
"price",
|
||||
"profit",
|
||||
})
|
||||
|
||||
|
||||
def _coerce_timeframe(timeframe: int | str) -> int:
|
||||
if isinstance(timeframe, int):
|
||||
@@ -419,325 +419,219 @@ class Mt5CliClient:
|
||||
return self._fetch(lambda c: c.market_book_get_as_df(symbol=symbol))
|
||||
|
||||
|
||||
def _create_cash_events_view(
|
||||
conn: sqlite3.Connection,
|
||||
deals_columns: set[str],
|
||||
) -> bool:
|
||||
"""Create the cash_events SQLite view derived from history_deals.
|
||||
def _resolve_incremental_settings(
|
||||
selected_datasets: set[Dataset],
|
||||
timeframes: Sequence[int | str] | None,
|
||||
flags: int | str,
|
||||
) -> tuple[list[int], int]:
|
||||
"""Resolve dataset-specific incremental update settings.
|
||||
|
||||
Returns:
|
||||
True if the view was created, False if required columns are missing.
|
||||
Tuple of resolved rate timeframes and tick copy flags.
|
||||
|
||||
Raises:
|
||||
ValueError: If timeframe or tick flag values are invalid.
|
||||
"""
|
||||
if "type" not in deals_columns:
|
||||
logger.warning("Skipping cash_events view: history_deals.type is missing")
|
||||
return False
|
||||
conn.execute("DROP VIEW IF EXISTS cash_events")
|
||||
conn.execute(
|
||||
"CREATE VIEW cash_events AS" # noqa: S608
|
||||
f" SELECT * FROM history_deals WHERE type NOT IN {_TRADE_DEAL_TYPES_SQL}",
|
||||
)
|
||||
return True
|
||||
resolved_timeframes: list[int] = []
|
||||
if Dataset.rates in selected_datasets:
|
||||
try:
|
||||
resolved_timeframes = resolve_history_timeframes(timeframes)
|
||||
except ValueError as exc:
|
||||
msg = str(exc)
|
||||
raise ValueError(msg) from exc
|
||||
resolved_tick_flags = 0
|
||||
if Dataset.ticks in selected_datasets:
|
||||
try:
|
||||
resolved_tick_flags = resolve_history_tick_flags(flags)
|
||||
except ValueError as exc:
|
||||
msg = str(exc)
|
||||
raise ValueError(msg) from exc
|
||||
return resolved_timeframes, resolved_tick_flags
|
||||
|
||||
|
||||
def _create_positions_reconstructed_view(
|
||||
conn: sqlite3.Connection,
|
||||
deals_columns: set[str],
|
||||
) -> bool:
|
||||
"""Create the positions_reconstructed SQLite view derived from history_deals.
|
||||
@dataclass(frozen=True)
|
||||
class _UpdateHistoryRequest:
|
||||
selected: set[Dataset]
|
||||
end: datetime
|
||||
fallback_start: datetime
|
||||
resolved_timeframes: list[int]
|
||||
resolved_tick_flags: int
|
||||
output_path: Path
|
||||
|
||||
|
||||
def _resolve_update_history_request(
|
||||
*,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str],
|
||||
datasets: set[Dataset] | None,
|
||||
timeframes: Sequence[int | str] | None,
|
||||
flags: int | str,
|
||||
lookback_hours: float,
|
||||
date_to: datetime | str | None,
|
||||
) -> _UpdateHistoryRequest | None:
|
||||
"""Validate and resolve incremental history update inputs.
|
||||
|
||||
Returns:
|
||||
True if the view was created, False if required columns are missing.
|
||||
Resolved request parameters, or None when no datasets are selected.
|
||||
|
||||
Raises:
|
||||
ValueError: If symbols are empty, lookback_hours is not positive, or
|
||||
timeframe/flag values are invalid.
|
||||
"""
|
||||
if not _POSITIONS_VIEW_REQUIRED_COLUMNS.issubset(deals_columns):
|
||||
missing = ", ".join(sorted(_POSITIONS_VIEW_REQUIRED_COLUMNS - deals_columns))
|
||||
logger.warning(
|
||||
"Skipping positions_reconstructed view: history_deals missing columns: %s",
|
||||
missing,
|
||||
)
|
||||
return False
|
||||
conn.execute("DROP VIEW IF EXISTS positions_reconstructed")
|
||||
conn.execute(
|
||||
"CREATE VIEW positions_reconstructed AS" # noqa: S608
|
||||
" SELECT"
|
||||
" position_id,"
|
||||
" symbol,"
|
||||
" MIN(CASE WHEN entry = 0 THEN time END) AS open_time,"
|
||||
" MAX(CASE WHEN entry IN (1, 2, 3) THEN time END) AS close_time,"
|
||||
" MIN(CASE WHEN entry = 0 THEN type END) AS direction,"
|
||||
" SUM(CASE WHEN entry = 0 THEN volume ELSE 0 END) AS volume_open,"
|
||||
" SUM(CASE WHEN entry IN (1, 3) THEN volume ELSE 0 END) AS volume_close,"
|
||||
" SUM(CASE WHEN entry = 2 THEN volume ELSE 0 END) AS volume_reversal,"
|
||||
" CASE"
|
||||
" WHEN SUM(CASE WHEN entry = 0 THEN volume ELSE 0 END) > 0"
|
||||
" THEN SUM(CASE WHEN entry = 0 THEN price * volume ELSE 0 END)"
|
||||
" / SUM(CASE WHEN entry = 0 THEN volume ELSE 0 END)"
|
||||
" END AS open_price,"
|
||||
" CASE"
|
||||
" WHEN SUM(CASE WHEN entry IN (1, 3) THEN volume ELSE 0 END) > 0"
|
||||
" THEN SUM(CASE WHEN entry IN (1, 3) THEN price * volume ELSE 0 END)"
|
||||
" / SUM(CASE WHEN entry IN (1, 3) THEN volume ELSE 0 END)"
|
||||
" END AS close_price,"
|
||||
" SUM(profit) AS total_profit,"
|
||||
" SUM(CASE WHEN entry = 2 THEN 1 ELSE 0 END) AS reversal_count,"
|
||||
" COUNT(*) AS deals_count"
|
||||
" FROM history_deals"
|
||||
f" WHERE type IN {_TRADE_DEAL_TYPES_SQL} AND position_id != 0"
|
||||
" GROUP BY position_id, symbol"
|
||||
" HAVING SUM(CASE WHEN entry IN (1, 3) THEN 1 ELSE 0 END) > 0",
|
||||
)
|
||||
return True
|
||||
if lookback_hours <= 0:
|
||||
msg = "lookback_hours must be positive."
|
||||
raise ValueError(msg)
|
||||
selected = resolve_history_datasets(datasets)
|
||||
if not selected:
|
||||
logger.info("Skipping SQLite history update: no datasets selected.")
|
||||
return None
|
||||
if not symbols:
|
||||
msg = "At least one symbol is required."
|
||||
raise ValueError(msg)
|
||||
|
||||
|
||||
def _write_frame_to_sqlite(
|
||||
conn: sqlite3.Connection,
|
||||
frame: pd.DataFrame,
|
||||
table_name: str,
|
||||
if_exists: IfExists,
|
||||
) -> bool:
|
||||
"""Write a non-empty-schema frame to SQLite.
|
||||
|
||||
Returns:
|
||||
True if a table was written, False if the frame had no columns.
|
||||
"""
|
||||
if len(frame.columns) == 0:
|
||||
logger.warning("Skipping %s: dataset returned no columns", table_name)
|
||||
return False
|
||||
frame.to_sql( # type: ignore[reportUnknownMemberType]
|
||||
table_name,
|
||||
conn,
|
||||
if_exists=if_exists.value,
|
||||
index=False,
|
||||
chunksize=50_000,
|
||||
method="multi",
|
||||
)
|
||||
return True
|
||||
|
||||
|
||||
def _create_collect_history_indexes(
|
||||
conn: sqlite3.Connection,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> None:
|
||||
"""Create useful indexes for collected history tables when present."""
|
||||
if {"symbol", "time"}.issubset(written_columns.get(Dataset.rates, set())):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_rates_symbol_time ON rates(symbol, time)",
|
||||
)
|
||||
if {"symbol", "time"}.issubset(written_columns.get(Dataset.ticks, set())):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_ticks_symbol_time ON ticks(symbol, time)",
|
||||
)
|
||||
if {"position_id", "symbol"}.issubset(
|
||||
written_columns.get(Dataset.history_deals, set())
|
||||
):
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_history_deals_position_symbol"
|
||||
" ON history_deals(position_id, symbol)",
|
||||
)
|
||||
|
||||
|
||||
def _record_written_columns(
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
dataset: Dataset,
|
||||
frame: pd.DataFrame,
|
||||
) -> None:
|
||||
"""Remember columns for datasets written during streaming collection."""
|
||||
columns = set(frame.columns)
|
||||
if dataset in written_columns:
|
||||
written_columns[dataset].update(columns)
|
||||
if date_to is not None:
|
||||
resolved_end = _coerce_datetime(date_to)
|
||||
else:
|
||||
written_columns[dataset] = columns
|
||||
|
||||
|
||||
def _write_streamed_frame(
|
||||
conn: sqlite3.Connection,
|
||||
frame: pd.DataFrame,
|
||||
dataset: Dataset,
|
||||
table_exists: bool,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Write one streamed dataset frame and track table state.
|
||||
|
||||
Returns:
|
||||
True if the dataset table exists after this write attempt.
|
||||
"""
|
||||
write_mode = IfExists.APPEND if table_exists else if_exists
|
||||
if _write_frame_to_sqlite(
|
||||
conn,
|
||||
frame,
|
||||
dataset.table_name,
|
||||
write_mode,
|
||||
):
|
||||
_record_written_columns(written_columns, dataset, frame)
|
||||
return True
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_rates_dataset(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
symbols: list[str],
|
||||
timeframe: int,
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Stream rates frames into SQLite.
|
||||
|
||||
Returns:
|
||||
True if the rates table was written.
|
||||
"""
|
||||
table_exists = False
|
||||
for sym in symbols:
|
||||
frame = client.copy_rates_range_as_df(
|
||||
symbol=sym,
|
||||
timeframe=timeframe,
|
||||
date_from=date_from,
|
||||
date_to=date_to,
|
||||
)
|
||||
frame.insert(0, "symbol", sym)
|
||||
frame.insert(1, "timeframe", timeframe)
|
||||
table_exists = _write_streamed_frame(
|
||||
conn,
|
||||
frame,
|
||||
Dataset.rates,
|
||||
table_exists,
|
||||
if_exists,
|
||||
written_columns,
|
||||
)
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_ticks_dataset(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
symbols: list[str],
|
||||
flags: int,
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Stream ticks frames into SQLite.
|
||||
|
||||
Returns:
|
||||
True if the ticks table was written.
|
||||
"""
|
||||
table_exists = False
|
||||
for sym in symbols:
|
||||
frame = client.copy_ticks_range_as_df(
|
||||
symbol=sym,
|
||||
date_from=date_from,
|
||||
date_to=date_to,
|
||||
flags=flags,
|
||||
)
|
||||
frame.insert(0, "symbol", sym)
|
||||
table_exists = _write_streamed_frame(
|
||||
conn,
|
||||
frame,
|
||||
Dataset.ticks,
|
||||
table_exists,
|
||||
if_exists,
|
||||
written_columns,
|
||||
)
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_history_dataset(
|
||||
conn: sqlite3.Connection,
|
||||
fetch: Callable[..., pd.DataFrame],
|
||||
dataset: Dataset,
|
||||
symbols: list[str],
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
written_columns: dict[Dataset, set[str]],
|
||||
) -> bool:
|
||||
"""Stream a history dataset into SQLite with exact symbol filtering.
|
||||
|
||||
Returns:
|
||||
True if the history table was written.
|
||||
"""
|
||||
table_exists = False
|
||||
for sym in symbols:
|
||||
frame = fetch(date_from=date_from, date_to=date_to, symbol=sym)
|
||||
if "symbol" in frame.columns:
|
||||
frame = frame[frame["symbol"] == sym]
|
||||
table_exists = _write_streamed_frame(
|
||||
conn,
|
||||
frame,
|
||||
dataset,
|
||||
table_exists,
|
||||
if_exists,
|
||||
written_columns,
|
||||
)
|
||||
return table_exists
|
||||
|
||||
|
||||
def _write_collected_datasets(
|
||||
conn: sqlite3.Connection,
|
||||
client: Mt5DataClient,
|
||||
symbols: list[str],
|
||||
datasets: set[Dataset],
|
||||
timeframe: int,
|
||||
flags: int,
|
||||
date_from: datetime,
|
||||
date_to: datetime,
|
||||
if_exists: IfExists,
|
||||
) -> tuple[set[Dataset], dict[Dataset, set[str]]]:
|
||||
"""Collect selected datasets and stream each symbol frame into SQLite.
|
||||
|
||||
Returns:
|
||||
Written datasets and their columns.
|
||||
"""
|
||||
written_columns: dict[Dataset, set[str]] = {}
|
||||
written_tables: set[Dataset] = set()
|
||||
if Dataset.rates in datasets and _write_rates_dataset(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
timeframe,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.rates)
|
||||
if Dataset.ticks in datasets and _write_ticks_dataset(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
resolved_end = datetime.now(UTC)
|
||||
end = resolved_end if resolved_end is not None else datetime.now(UTC)
|
||||
fallback_start = end - timedelta(hours=lookback_hours)
|
||||
resolved_timeframes, resolved_tick_flags = _resolve_incremental_settings(
|
||||
selected,
|
||||
timeframes,
|
||||
flags,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.ticks)
|
||||
if Dataset.history_orders in datasets and _write_history_dataset(
|
||||
conn,
|
||||
client.history_orders_get_as_df,
|
||||
Dataset.history_orders,
|
||||
symbols,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.history_orders)
|
||||
if Dataset.history_deals in datasets and _write_history_dataset(
|
||||
conn,
|
||||
client.history_deals_get_as_df,
|
||||
Dataset.history_deals,
|
||||
symbols,
|
||||
date_from,
|
||||
date_to,
|
||||
if_exists,
|
||||
written_columns,
|
||||
):
|
||||
written_tables.add(Dataset.history_deals)
|
||||
return written_tables, written_columns
|
||||
)
|
||||
return _UpdateHistoryRequest(
|
||||
selected=selected,
|
||||
end=end,
|
||||
fallback_start=fallback_start,
|
||||
resolved_timeframes=resolved_timeframes,
|
||||
resolved_tick_flags=resolved_tick_flags,
|
||||
output_path=Path(output),
|
||||
)
|
||||
|
||||
|
||||
def update_history( # noqa: PLR0913
|
||||
*,
|
||||
client: Mt5DataClient,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str],
|
||||
datasets: set[Dataset] | None = None,
|
||||
timeframes: Sequence[int | str] | None = None,
|
||||
flags: int | str = "ALL",
|
||||
lookback_hours: float = 24.0,
|
||||
date_to: datetime | str | None = None,
|
||||
deduplicate: bool = True,
|
||||
create_rate_views: bool = True,
|
||||
with_views: bool = False,
|
||||
include_account_events: bool = True,
|
||||
) -> None:
|
||||
"""Incrementally append MT5 history into a SQLite database.
|
||||
|
||||
Uses an already-connected ``Mt5DataClient`` and does not create or close
|
||||
the MT5 connection. For first-time tables, data is fetched from
|
||||
``date_to - lookback_hours``. Subsequent runs resume from existing
|
||||
``MAX(time)`` per symbol (and timeframe for rates); when
|
||||
``include_account_events=True``, account-level deals use a separate cursor
|
||||
over ``type NOT IN (0, 1)`` / empty-symbol rows.
|
||||
|
||||
Args:
|
||||
client: Connected MT5 data client.
|
||||
output: SQLite database path.
|
||||
symbols: Symbols to update.
|
||||
datasets: Datasets to include (defaults to all).
|
||||
timeframes: Rate timeframes to update (defaults to all fixed MT5
|
||||
timeframes when None).
|
||||
flags: Tick copy flags as integer or name (e.g. ``ALL``).
|
||||
lookback_hours: First-run lookback when a table has no prior rows.
|
||||
date_to: Optional update end datetime. Defaults to now (UTC).
|
||||
deduplicate: Remove duplicate rows after append, keeping latest ROWID.
|
||||
create_rate_views: Create ``rate_<symbol>__<timeframe>`` views.
|
||||
with_views: Create ``cash_events`` and ``positions_reconstructed`` views.
|
||||
include_account_events: Include account-level cash events in
|
||||
``history_deals`` when True.
|
||||
"""
|
||||
request = _resolve_update_history_request(
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
datasets=datasets,
|
||||
timeframes=timeframes,
|
||||
flags=flags,
|
||||
lookback_hours=lookback_hours,
|
||||
date_to=date_to,
|
||||
)
|
||||
if request is None:
|
||||
return
|
||||
logger.info(
|
||||
"Updating history in SQLite: symbols=%s, datasets=%s, path=%s",
|
||||
list(symbols),
|
||||
sorted(dataset.value for dataset in request.selected),
|
||||
request.output_path,
|
||||
)
|
||||
with sqlite3.connect(request.output_path) as conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
write_incremental_datasets(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
request.selected,
|
||||
request.resolved_timeframes,
|
||||
request.resolved_tick_flags,
|
||||
request.fallback_start,
|
||||
request.end,
|
||||
deduplicate=deduplicate,
|
||||
create_rate_views=create_rate_views,
|
||||
with_views=with_views,
|
||||
include_account_events=include_account_events,
|
||||
)
|
||||
|
||||
|
||||
def update_history_with_config( # noqa: PLR0913
|
||||
*,
|
||||
output: Path | str,
|
||||
symbols: Sequence[str],
|
||||
config: Mt5Config | None = None,
|
||||
datasets: set[Dataset] | None = None,
|
||||
timeframes: Sequence[int | str] | None = None,
|
||||
flags: int | str = "ALL",
|
||||
lookback_hours: float = 24.0,
|
||||
date_to: datetime | str | None = None,
|
||||
deduplicate: bool = True,
|
||||
create_rate_views: bool = True,
|
||||
with_views: bool = False,
|
||||
include_account_events: bool = True,
|
||||
) -> None:
|
||||
"""Incrementally append MT5 history, opening and closing the MT5 connection.
|
||||
|
||||
Convenience wrapper around :func:`update_history` for standalone use.
|
||||
"""
|
||||
request = _resolve_update_history_request(
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
datasets=datasets,
|
||||
timeframes=timeframes,
|
||||
flags=flags,
|
||||
lookback_hours=lookback_hours,
|
||||
date_to=date_to,
|
||||
)
|
||||
if request is None:
|
||||
return
|
||||
mt5_config = config or build_config()
|
||||
with _connected_client(mt5_config) as client:
|
||||
update_history(
|
||||
client=client,
|
||||
output=output,
|
||||
symbols=symbols,
|
||||
datasets=datasets,
|
||||
timeframes=timeframes,
|
||||
flags=flags,
|
||||
lookback_hours=lookback_hours,
|
||||
date_to=date_to,
|
||||
deduplicate=deduplicate,
|
||||
create_rate_views=create_rate_views,
|
||||
with_views=with_views,
|
||||
include_account_events=include_account_events,
|
||||
)
|
||||
|
||||
|
||||
def collect_history(
|
||||
@@ -776,7 +670,7 @@ def collect_history(
|
||||
with _connected_client(mt5_config) as client, sqlite3.connect(output) as conn:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA synchronous=NORMAL")
|
||||
written_tables, written_columns = _write_collected_datasets(
|
||||
written_tables, written_columns = write_collected_datasets(
|
||||
conn,
|
||||
client,
|
||||
symbols,
|
||||
@@ -787,10 +681,10 @@ def collect_history(
|
||||
end,
|
||||
if_exists,
|
||||
)
|
||||
_create_collect_history_indexes(conn, written_columns)
|
||||
create_history_indexes(conn, written_columns)
|
||||
if with_views and Dataset.history_deals in written_tables:
|
||||
_create_cash_events_view(conn, written_columns[Dataset.history_deals])
|
||||
_create_positions_reconstructed_view(
|
||||
create_cash_events_view(conn, written_columns[Dataset.history_deals])
|
||||
create_positions_reconstructed_view(
|
||||
conn,
|
||||
written_columns[Dataset.history_deals],
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user