Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ad9e513253 |
+2
-2
@@ -133,8 +133,8 @@ The `update_history` SDK path uses the same base tables and optional
|
|||||||
### Rate view resolution
|
### Rate view resolution
|
||||||
|
|
||||||
Downstream tools can resolve mt5cli-managed compatibility view names from an
|
Downstream tools can resolve mt5cli-managed compatibility view names from an
|
||||||
existing SQLite history database without creating files or guessing legacy
|
existing SQLite history database without creating files or guessing naming
|
||||||
naming schemes:
|
schemes:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|||||||
+40
-9
@@ -21,7 +21,7 @@ from .utils import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from collections.abc import Callable, Sequence
|
from collections.abc import Callable, Mapping, Sequence
|
||||||
|
|
||||||
from pdmt5 import Mt5DataClient
|
from pdmt5 import Mt5DataClient
|
||||||
|
|
||||||
@@ -989,7 +989,20 @@ def drop_duplicates_in_table(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
DedupScope = tuple[str, tuple[object, ...]]
|
@dataclass(frozen=True)
|
||||||
|
class DedupScope:
|
||||||
|
"""Scoped deduplication predicate and the columns it references.
|
||||||
|
|
||||||
|
Attributes:
|
||||||
|
where: SQL predicate appended to the duplicate-removal query.
|
||||||
|
params: Parameters bound to the scope predicate.
|
||||||
|
required_columns: Columns that must be present in the written table for
|
||||||
|
the scope to run.
|
||||||
|
"""
|
||||||
|
|
||||||
|
where: str
|
||||||
|
params: tuple[object, ...]
|
||||||
|
required_columns: frozenset[str]
|
||||||
|
|
||||||
|
|
||||||
def _record_dedup_scope(
|
def _record_dedup_scope(
|
||||||
@@ -997,17 +1010,25 @@ def _record_dedup_scope(
|
|||||||
dataset: Dataset,
|
dataset: Dataset,
|
||||||
scope_where: str,
|
scope_where: str,
|
||||||
scope_params: tuple[object, ...],
|
scope_params: tuple[object, ...],
|
||||||
|
required_columns: frozenset[str],
|
||||||
) -> None:
|
) -> None:
|
||||||
dedup_scopes.setdefault(dataset, []).append((scope_where, scope_params))
|
dedup_scopes.setdefault(dataset, []).append(
|
||||||
|
DedupScope(scope_where, scope_params, required_columns),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def deduplicate_history_tables(
|
def deduplicate_history_tables(
|
||||||
conn: sqlite3.Connection,
|
conn: sqlite3.Connection,
|
||||||
written_columns: dict[Dataset, set[str]],
|
written_columns: dict[Dataset, set[str]],
|
||||||
written_tables: set[Dataset],
|
written_tables: set[Dataset],
|
||||||
dedup_scopes: dict[Dataset, list[DedupScope]] | None = None,
|
dedup_scopes: Mapping[Dataset, Sequence[DedupScope]] | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Deduplicate appended history tables by stable identifiers."""
|
"""Deduplicate appended history tables by stable identifiers.
|
||||||
|
|
||||||
|
Scopes whose required columns are not present in the written table are
|
||||||
|
skipped. If all scopes for a dataset are skipped, the table receives one
|
||||||
|
unscoped deduplication pass instead.
|
||||||
|
"""
|
||||||
cursor = conn.cursor()
|
cursor = conn.cursor()
|
||||||
for dataset in written_tables:
|
for dataset in written_tables:
|
||||||
columns = written_columns.get(dataset, set())
|
columns = written_columns.get(dataset, set())
|
||||||
@@ -1026,16 +1047,19 @@ def deduplicate_history_tables(
|
|||||||
table,
|
table,
|
||||||
)
|
)
|
||||||
continue
|
continue
|
||||||
scopes = dedup_scopes.get(dataset, []) if dedup_scopes else []
|
raw_scopes: Sequence[DedupScope] = (
|
||||||
|
dedup_scopes.get(dataset, ()) if dedup_scopes else ()
|
||||||
|
)
|
||||||
|
scopes = [scope for scope in raw_scopes if scope.required_columns <= columns]
|
||||||
if scopes:
|
if scopes:
|
||||||
for scope_where, scope_params in scopes:
|
for scope in scopes:
|
||||||
drop_duplicates_in_table(
|
drop_duplicates_in_table(
|
||||||
cursor,
|
cursor,
|
||||||
table,
|
table,
|
||||||
list(keys),
|
list(keys),
|
||||||
keep="last",
|
keep="last",
|
||||||
scope_where=scope_where,
|
scope_where=scope.where,
|
||||||
scope_params=scope_params,
|
scope_params=scope.params,
|
||||||
)
|
)
|
||||||
continue
|
continue
|
||||||
drop_duplicates_in_table(cursor, table, list(keys), keep="last")
|
drop_duplicates_in_table(cursor, table, list(keys), keep="last")
|
||||||
@@ -1402,6 +1426,7 @@ def _write_incremental_rates(
|
|||||||
Dataset.rates,
|
Dataset.rates,
|
||||||
"symbol = ? AND timeframe = ? AND time >= ?",
|
"symbol = ? AND timeframe = ? AND time >= ?",
|
||||||
(symbol, timeframe, start_date),
|
(symbol, timeframe, start_date),
|
||||||
|
frozenset({"symbol", "timeframe", "time"}),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -1440,6 +1465,7 @@ def _write_incremental_ticks(
|
|||||||
Dataset.ticks,
|
Dataset.ticks,
|
||||||
"symbol = ? AND time >= ?",
|
"symbol = ? AND time >= ?",
|
||||||
(symbol, start_date),
|
(symbol, start_date),
|
||||||
|
frozenset({"symbol", "time"}),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -1478,6 +1504,7 @@ def _write_incremental_history_orders(
|
|||||||
Dataset.history_orders,
|
Dataset.history_orders,
|
||||||
"symbol = ? AND time >= ?",
|
"symbol = ? AND time >= ?",
|
||||||
(symbol, start_date),
|
(symbol, start_date),
|
||||||
|
frozenset({"symbol", "time"}),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -1531,6 +1558,7 @@ def _write_incremental_history_deals(
|
|||||||
Dataset.history_deals,
|
Dataset.history_deals,
|
||||||
"symbol = ? AND time >= ?",
|
"symbol = ? AND time >= ?",
|
||||||
(symbol, start_by_symbol[symbol, None]),
|
(symbol, start_by_symbol[symbol, None]),
|
||||||
|
frozenset({"symbol", "time"}),
|
||||||
)
|
)
|
||||||
if "type" in columns:
|
if "type" in columns:
|
||||||
_record_dedup_scope(
|
_record_dedup_scope(
|
||||||
@@ -1538,6 +1566,7 @@ def _write_incremental_history_deals(
|
|||||||
Dataset.history_deals,
|
Dataset.history_deals,
|
||||||
f"type NOT IN {_TRADE_DEAL_TYPES_SQL} AND time >= ?",
|
f"type NOT IN {_TRADE_DEAL_TYPES_SQL} AND time >= ?",
|
||||||
(account_event_start,),
|
(account_event_start,),
|
||||||
|
frozenset({"type", "time"}),
|
||||||
)
|
)
|
||||||
if "type" not in columns and "symbol" in columns:
|
if "type" not in columns and "symbol" in columns:
|
||||||
_record_dedup_scope(
|
_record_dedup_scope(
|
||||||
@@ -1545,6 +1574,7 @@ def _write_incremental_history_deals(
|
|||||||
Dataset.history_deals,
|
Dataset.history_deals,
|
||||||
"(symbol IS NULL OR symbol = '') AND time >= ?",
|
"(symbol IS NULL OR symbol = '') AND time >= ?",
|
||||||
(account_event_start,),
|
(account_event_start,),
|
||||||
|
frozenset({"symbol", "time"}),
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
start_by_symbol = load_incremental_start_datetimes(
|
start_by_symbol = load_incremental_start_datetimes(
|
||||||
@@ -1572,6 +1602,7 @@ def _write_incremental_history_deals(
|
|||||||
Dataset.history_deals,
|
Dataset.history_deals,
|
||||||
"symbol = ? AND time >= ?",
|
"symbol = ? AND time >= ?",
|
||||||
(symbol, start_date),
|
(symbol, start_date),
|
||||||
|
frozenset({"symbol", "time"}),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
[project]
|
[project]
|
||||||
name = "mt5cli"
|
name = "mt5cli"
|
||||||
version = "0.5.1"
|
version = "0.5.2"
|
||||||
description = "Command-line tool for MetaTrader 5"
|
description = "Command-line tool for MetaTrader 5"
|
||||||
authors = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
authors = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||||
maintainers = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
maintainers = [{name = "dceoy", email = "dceoy@users.noreply.github.com"}]
|
||||||
|
|||||||
+136
-3
@@ -18,6 +18,7 @@ if TYPE_CHECKING:
|
|||||||
from mt5cli import history
|
from mt5cli import history
|
||||||
from mt5cli.history import (
|
from mt5cli.history import (
|
||||||
DEFAULT_HISTORY_TIMEFRAMES,
|
DEFAULT_HISTORY_TIMEFRAMES,
|
||||||
|
DedupScope,
|
||||||
RateTarget,
|
RateTarget,
|
||||||
append_dataframe,
|
append_dataframe,
|
||||||
augment_written_columns_from_sqlite,
|
augment_written_columns_from_sqlite,
|
||||||
@@ -618,7 +619,7 @@ class TestIncrementalStart:
|
|||||||
) -> None:
|
) -> None:
|
||||||
"""Test rates tables without timeframe fail fast during incremental resume."""
|
"""Test rates tables without timeframe fail fast during incremental resume."""
|
||||||
fallback = datetime(2024, 1, 1, tzinfo=UTC)
|
fallback = datetime(2024, 1, 1, tzinfo=UTC)
|
||||||
with sqlite3.connect(tmp_path / "legacy-rates.db") as conn:
|
with sqlite3.connect(tmp_path / "rates-without-timeframe.db") as conn:
|
||||||
conn.execute("CREATE TABLE rates(symbol TEXT, time TEXT, open REAL)")
|
conn.execute("CREATE TABLE rates(symbol TEXT, time TEXT, open REAL)")
|
||||||
conn.execute(
|
conn.execute(
|
||||||
"INSERT INTO rates(symbol, time, open) VALUES (?, ?, ?)",
|
"INSERT INTO rates(symbol, time, open) VALUES (?, ?, ?)",
|
||||||
@@ -887,9 +888,10 @@ class TestDeduplication:
|
|||||||
{Dataset.rates},
|
{Dataset.rates},
|
||||||
{
|
{
|
||||||
Dataset.rates: [
|
Dataset.rates: [
|
||||||
(
|
DedupScope(
|
||||||
"symbol = ? AND timeframe = ? AND time >= ?",
|
"symbol = ? AND timeframe = ? AND time >= ?",
|
||||||
("EURUSD", 1, boundary),
|
("EURUSD", 1, boundary),
|
||||||
|
frozenset({"symbol", "timeframe", "time"}),
|
||||||
),
|
),
|
||||||
],
|
],
|
||||||
},
|
},
|
||||||
@@ -902,6 +904,89 @@ class TestDeduplication:
|
|||||||
("2024-01-02T00:00:00+00:00", 9.9),
|
("2024-01-02T00:00:00+00:00", 9.9),
|
||||||
]
|
]
|
||||||
|
|
||||||
|
def test_unusable_scope_falls_back_to_table_dedup(self, tmp_path: Path) -> None:
|
||||||
|
"""Test scopes with missing columns do not break stable-key dedup."""
|
||||||
|
boundary = datetime(2024, 1, 1, tzinfo=UTC)
|
||||||
|
with sqlite3.connect(tmp_path / "orders-without-time.db") as conn:
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE history_orders("
|
||||||
|
" ticket INTEGER, symbol TEXT, time_setup TEXT, type INTEGER)",
|
||||||
|
)
|
||||||
|
conn.executemany(
|
||||||
|
"INSERT INTO history_orders(ticket, symbol, time_setup, type)"
|
||||||
|
" VALUES (?, ?, ?, ?)",
|
||||||
|
[
|
||||||
|
(1, "EURUSD", "2024-01-01T00:00:00+00:00", 0),
|
||||||
|
(1, "EURUSD", "2024-01-01T00:00:01+00:00", 1),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
deduplicate_history_tables(
|
||||||
|
conn,
|
||||||
|
{Dataset.history_orders: {"ticket", "symbol", "time_setup", "type"}},
|
||||||
|
{Dataset.history_orders},
|
||||||
|
{
|
||||||
|
Dataset.history_orders: [
|
||||||
|
DedupScope(
|
||||||
|
"symbol = ? AND time >= ?",
|
||||||
|
("EURUSD", boundary),
|
||||||
|
frozenset({"symbol", "time"}),
|
||||||
|
),
|
||||||
|
],
|
||||||
|
},
|
||||||
|
)
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT ticket, time_setup, type FROM history_orders",
|
||||||
|
).fetchall()
|
||||||
|
assert rows == [(1, "2024-01-01T00:00:01+00:00", 1)]
|
||||||
|
|
||||||
|
def test_partially_unusable_scopes_only_run_usable_scopes(
|
||||||
|
self,
|
||||||
|
tmp_path: Path,
|
||||||
|
) -> None:
|
||||||
|
"""Test mixed scope filtering skips only scopes with missing columns."""
|
||||||
|
boundary = datetime(2024, 1, 2, tzinfo=UTC)
|
||||||
|
with sqlite3.connect(tmp_path / "partial-scope-filter.db") as conn:
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE rates("
|
||||||
|
" symbol TEXT, timeframe INTEGER, time TEXT, open REAL)",
|
||||||
|
)
|
||||||
|
conn.executemany(
|
||||||
|
"INSERT INTO rates(symbol, timeframe, time, open) VALUES (?, ?, ?, ?)",
|
||||||
|
[
|
||||||
|
("EURUSD", 1, "2024-01-02T00:00:00+00:00", 2.0),
|
||||||
|
("EURUSD", 1, "2024-01-02T00:00:00+00:00", 9.9),
|
||||||
|
("USDJPY", 1, "2024-01-02T00:00:00+00:00", 100.0),
|
||||||
|
("USDJPY", 1, "2024-01-02T00:00:00+00:00", 101.0),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
deduplicate_history_tables(
|
||||||
|
conn,
|
||||||
|
{Dataset.rates: {"symbol", "timeframe", "time", "open"}},
|
||||||
|
{Dataset.rates},
|
||||||
|
{
|
||||||
|
Dataset.rates: [
|
||||||
|
DedupScope(
|
||||||
|
"symbol = ? AND timeframe = ? AND time >= ?",
|
||||||
|
("EURUSD", 1, boundary),
|
||||||
|
frozenset({"symbol", "timeframe", "time"}),
|
||||||
|
),
|
||||||
|
DedupScope(
|
||||||
|
"symbol = ? AND timeframe = ? AND broker = ?",
|
||||||
|
("USDJPY", 1, "demo"),
|
||||||
|
frozenset({"symbol", "timeframe", "broker"}),
|
||||||
|
),
|
||||||
|
],
|
||||||
|
},
|
||||||
|
)
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT symbol, open FROM rates ORDER BY symbol, open",
|
||||||
|
).fetchall()
|
||||||
|
assert rows == [
|
||||||
|
("EURUSD", 9.9),
|
||||||
|
("USDJPY", 100.0),
|
||||||
|
("USDJPY", 101.0),
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
class TestRateCompatibilityViews:
|
class TestRateCompatibilityViews:
|
||||||
"""Tests for rate compatibility view creation."""
|
"""Tests for rate compatibility view creation."""
|
||||||
@@ -1362,6 +1447,54 @@ class TestIncrementalIntegration:
|
|||||||
"rate_EURUSD_M1__1",
|
"rate_EURUSD_M1__1",
|
||||||
}
|
}
|
||||||
|
|
||||||
|
def test_incremental_orders_without_time_deduplicate_by_ticket(
|
||||||
|
self,
|
||||||
|
tmp_path: Path,
|
||||||
|
caplog: pytest.LogCaptureFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Test incremental history_orders without time deduplicate safely."""
|
||||||
|
|
||||||
|
def history_orders_get_as_df(**kwargs: object) -> pd.DataFrame:
|
||||||
|
if kwargs["symbol"] == "GBPUSD":
|
||||||
|
return pd.DataFrame()
|
||||||
|
return pd.DataFrame({
|
||||||
|
"ticket": [1, 1],
|
||||||
|
"symbol": ["EURUSD", "EURUSD"],
|
||||||
|
"time_setup": [
|
||||||
|
"2024-01-01T00:00:00+00:00",
|
||||||
|
"2024-01-01T00:00:01+00:00",
|
||||||
|
],
|
||||||
|
"type": [0, 1],
|
||||||
|
})
|
||||||
|
|
||||||
|
client = MagicMock()
|
||||||
|
client.history_orders_get_as_df.side_effect = history_orders_get_as_df
|
||||||
|
start = datetime(2024, 1, 1, tzinfo=UTC)
|
||||||
|
end = datetime(2024, 1, 2, tzinfo=UTC)
|
||||||
|
with (
|
||||||
|
sqlite3.connect(tmp_path / "incremental-orders-without-time.db") as conn,
|
||||||
|
caplog.at_level(logging.WARNING, logger="mt5cli.history"),
|
||||||
|
):
|
||||||
|
write_incremental_datasets(
|
||||||
|
conn,
|
||||||
|
client,
|
||||||
|
["EURUSD", "GBPUSD"],
|
||||||
|
{Dataset.history_orders},
|
||||||
|
[],
|
||||||
|
0,
|
||||||
|
start,
|
||||||
|
end,
|
||||||
|
deduplicate=True,
|
||||||
|
create_rate_views=False,
|
||||||
|
with_views=False,
|
||||||
|
include_account_events=False,
|
||||||
|
)
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT ticket, time_setup, type FROM history_orders",
|
||||||
|
).fetchall()
|
||||||
|
assert rows == [(1, "2024-01-01T00:00:01+00:00", 1)]
|
||||||
|
assert "Skipping history_orders: dataset returned no columns" in caplog.text
|
||||||
|
|
||||||
def test_write_collected_datasets_and_edge_branches(
|
def test_write_collected_datasets_and_edge_branches(
|
||||||
self,
|
self,
|
||||||
tmp_path: Path,
|
tmp_path: Path,
|
||||||
@@ -1737,7 +1870,7 @@ class TestIncrementalHistoryDeals:
|
|||||||
})
|
})
|
||||||
start = datetime(2024, 1, 1, tzinfo=UTC)
|
start = datetime(2024, 1, 1, tzinfo=UTC)
|
||||||
end = datetime(2024, 1, 3, tzinfo=UTC)
|
end = datetime(2024, 1, 3, tzinfo=UTC)
|
||||||
with sqlite3.connect(tmp_path / "legacy-deals.db") as conn:
|
with sqlite3.connect(tmp_path / "deals-without-type.db") as conn:
|
||||||
conn.execute(
|
conn.execute(
|
||||||
"CREATE TABLE history_deals( ticket INTEGER, symbol TEXT, time TEXT)",
|
"CREATE TABLE history_deals( ticket INTEGER, symbol TEXT, time TEXT)",
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user