From 5a9e5eff698235f92f58b99384edbf64a5a4ceb1 Mon Sep 17 00:00:00 2001 From: DoRmAmMu1997 Date: Tue, 25 Aug 2026 00:28:24 +0530 Subject: [PATCH 1/2] fix(DATA-004): stop re-downloading stocks that listed after the window opened Both cache-coverage checks require first_date <= requested_start, where requested_start is today minus ten years. A stock that listed after that date can never satisfy it, because DhanHQ has nothing earlier to give. 200 of 577 cached symbols were in that state (DMART listed 2017-03-21, RBLBANK 2016-08-31, LTTS, COHANCE, ~196 more): every prefetch and every scan re-downloaded their full history, wrote the same short frame back, and did it again next time. The two cases look identical from the cached file alone and must not be conflated: an interrupted prefetch leaves a partial file while the vendor DOES have the missing years and must be refetched; a later listing means the vendor has nothing earlier and refetching is waste forever. Fixed by recording what the vendor actually served, following the sidecar precedent DATA-002 established with .checked and .repaired. A new .firstbar marker stores the probe's requested_from, the earliest_available bar that came back, and recorded_on. A cache is then treated as reaching back far enough when it either literally covers the requested start or already begins at the vendor's earliest known bar. Three properties keep the marker from ever hiding real missing history: - it only counts when the recorded probe reached at least as far back as the current request, so a five-year probe cannot suppress a ten-year refetch; - it expires after VENDOR_EARLIEST_RECHECK_DAYS (30), because vendors do occasionally backfill history; - any read or parse failure returns None and the strict rule applies, so a corrupt marker can only cost an extra request. Suppressing the backfill deliberately falls through to the normal freshness and incremental logic rather than returning early: a later listing still needs its daily top-up, and the point is to stop the pointless ten-year refetch, not to freeze the symbol. Measured on a copy of the real 577-file cache, nifty_500 (500 rows): prefetch pass 1 500 requests, 176 full-window backfills (learns) prefetch pass 2 176 requests, 0 full-window backfills prefetch pass 3 0 requests, 0 full-window backfills scan path before: 324 hits / 176 misses (176 full re-downloads) after: 500 hits / 0 misses, 0 files rewritten Stacked on DATA-003 (PR #112), which introduced _cache_covers_range and is not yet merged; this branch targets that one. Co-Authored-By: Claude Opus 5 --- backend/daily_data_loader.py | 192 ++++++++++- .../test_daily_data_loader_vendor_earliest.py | 311 ++++++++++++++++++ 2 files changed, 494 insertions(+), 9 deletions(-) create mode 100644 tests/test_daily_data_loader_vendor_earliest.py diff --git a/backend/daily_data_loader.py b/backend/daily_data_loader.py index 1b44a97..1e3475c 100644 --- a/backend/daily_data_loader.py +++ b/backend/daily_data_loader.py @@ -8,6 +8,7 @@ from __future__ import annotations import concurrent.futures +import json import logging import re import threading @@ -49,6 +50,13 @@ # ``history_start_date``) in one place stops the three callers from drifting apart. DEFAULT_HISTORY_YEARS_BACK = 10 +# How long we trust a recorded "the vendor has nothing earlier than this" answer +# before probing again. Vendors do occasionally backfill history, so the belief +# expires rather than becoming permanent; a month keeps the cost negligible +# (one request per affected symbol per month) while never hiding real data for +# long. Mirrors the DATA-002 ``REPAIR_RETRY_AFTER_DAYS`` cooldown precedent. +VENDOR_EARLIEST_RECHECK_DAYS = 30 + def history_start_date( years_back: int = DEFAULT_HISTORY_YEARS_BACK, today: date | None = None @@ -109,6 +117,7 @@ def _cache_covers_range( last_date: date | None, requested_start: date, requested_end: date, + vendor_earliest: date | None = None, ) -> bool: """Return True when a cached range is good enough to answer a request. @@ -130,7 +139,35 @@ def _cache_covers_range( if first_date is None or last_date is None: return False tolerated_end = requested_end - timedelta(days=STALE_LATEST_TOLERANCE_DAYS) - return first_date <= requested_start and last_date >= tolerated_end + if last_date < tolerated_end: + return False + return _cache_reaches_back_far_enough(first_date, requested_start, vendor_earliest) + + +def _cache_reaches_back_far_enough( + first_date: date, + requested_start: date, + vendor_earliest: date | None, +) -> bool: + """Return True when nothing earlier is missing that could still be fetched. + + Normally that means the cache literally reaches ``requested_start``. But a + stock that listed *after* that date can never satisfy it — DhanHQ has nothing + earlier to give — and demanding it made 200 of 577 symbols a permanent cache + miss, re-downloading their full history on every prefetch and every scan + (DATA-004). + + ``vendor_earliest`` is the earliest bar the vendor actually served for a probe + that reached at least as far back as this request (see + ``DailyDataLoader._vendor_earliest_for``). When the cache already starts at or + before that bar, it is as complete as it can ever be, so asking again is pure + waste. ``None`` means we have no such evidence and the strict rule applies — + which keeps an interrupted prefetch's partial file being refetched, since + there the vendor genuinely does have the missing years. + """ + if first_date <= requested_start: + return True + return vendor_earliest is not None and first_date <= vendor_earliest def _date_bounds(candles: pd.DataFrame) -> tuple[date | None, date | None]: @@ -333,6 +370,106 @@ def _write_checked_through( except OSError: logger.warning("Could not write daily-cache checked marker for %s", symbol) + def first_bar_path(self, symbol: str, security_id: str | int) -> Path: + """Return the sidecar recording how far back the vendor's history goes. + + A third marker alongside ``.checked`` (an empty tail) and ``.repaired`` + (a repair cooldown), following the same pattern: remember an answer the + vendor already gave so we do not pay for the identical request forever. + """ + return self.cache_path(symbol, security_id).with_suffix(".firstbar") + + def _vendor_earliest_for( + self, symbol: str, security_id: str | int, requested_start: date, today: date + ) -> date | None: + """The vendor's earliest bar, when we have evidence that answers this request. + + Returns ``None`` — meaning "no evidence, apply the strict rule" — unless + all three hold: + + - a marker exists and parses; + - it was recorded within ``VENDOR_EARLIEST_RECHECK_DAYS`` (vendors do + backfill occasionally, so the belief expires); + - the recorded probe reached **at least as far back** as this request. + Learning that nothing exists before 2021 when you only asked from 2021 + says nothing about 2016, so a shallower probe must not suppress a + deeper refetch. + + Any read or parse problem returns ``None``, so a corrupt marker can only + ever cost an extra request — never hide history that really is missing. + """ + path = self.first_bar_path(symbol, security_id) + if not path.exists(): + return None + try: + payload = json.loads(path.read_text(encoding="utf-8")) + requested_from = _coerce_date(str(payload["requested_from"])) + earliest_available = _coerce_date(str(payload["earliest_available"])) + recorded_on = _coerce_date(str(payload["recorded_on"])) + except (OSError, ValueError, KeyError, TypeError, json.JSONDecodeError): + return None + if (today - recorded_on).days >= VENDOR_EARLIEST_RECHECK_DAYS: + return None + if requested_from > requested_start: + return None + return earliest_available + + def _write_vendor_earliest( + self, + symbol: str, + security_id: str | int, + *, + requested_from: date, + earliest_available: date, + recorded_on: date, + ) -> None: + """Persist what the vendor served, so the next pass need not ask again.""" + try: + self.first_bar_path(symbol, security_id).write_text( + json.dumps( + { + "requested_from": requested_from.isoformat(), + "earliest_available": earliest_available.isoformat(), + "recorded_on": recorded_on.isoformat(), + } + ), + encoding="utf-8", + ) + except OSError: + # The marker is an optimisation, never a correctness requirement. + logger.warning("Could not write daily-cache first-bar marker for %s", symbol) + + def _record_vendor_earliest( + self, + symbol: str, + security_id: str | int, + *, + requested_from: date | datetime | str, + candles: pd.DataFrame, + today: date, + ) -> None: + """Record the vendor's earliest bar when it fell short of what we asked for. + + Called after any full-window download. A response that *does* reach the + requested start teaches us nothing worth storing, and an empty response is + no evidence at all, so both are skipped. + """ + if candles.empty: + return + first_date, _last_date = _date_bounds(candles) + if first_date is None: + return + start = _coerce_date(requested_from) + if first_date <= start: + return + self._write_vendor_earliest( + symbol, + security_id, + requested_from=start, + earliest_available=first_date, + recorded_on=today, + ) + def read_cached_history(self, symbol: str, security_id: str | int) -> pd.DataFrame: """Return the cached daily candles for one stock; empty DataFrame if missing. @@ -411,7 +548,12 @@ def get_daily_history( first_date, last_date = _date_bounds(cached) requested_start = _coerce_date(start_date) requested_end = _coerce_date(end_date) - if _cache_covers_range(first_date, last_date, requested_start, requested_end): + vendor_earliest = self._vendor_earliest_for( + symbol, security_id, requested_start, requested_end + ) + if _cache_covers_range( + first_date, last_date, requested_start, requested_end, vendor_earliest + ): if cached is None: # The footer is only an advisory index. The file can be # replaced after the metadata read, and a valid footer @@ -419,7 +561,11 @@ def get_daily_history( cached = pd.read_parquet(path) actual_first, actual_last = _date_bounds(cached) if _cache_covers_range( - actual_first, actual_last, requested_start, requested_end + actual_first, + actual_last, + requested_start, + requested_end, + vendor_earliest, ): return self._slice_to_range(cached, start_date, end_date), True @@ -435,6 +581,13 @@ def get_daily_history( if not candles.empty: path.parent.mkdir(parents=True, exist_ok=True) candles.to_parquet(path, index=False) + self._record_vendor_earliest( + symbol, + security_id, + requested_from=start_date, + candles=candles, + today=_coerce_date(end_date), + ) return self._slice_to_range(candles, start_date, end_date), False def fetch_window( @@ -516,6 +669,9 @@ def ensure_daily_history( if not candles.empty: path.parent.mkdir(parents=True, exist_ok=True) candles.to_parquet(path, index=False) + self._record_vendor_earliest( + symbol, security_id, requested_from=start, candles=candles, today=today + ) return candles, "fresh_download" cached = pd.read_parquet(path) @@ -531,6 +687,9 @@ def ensure_daily_history( ) if not candles.empty: candles.to_parquet(path, index=False) + self._record_vendor_earliest( + symbol, security_id, requested_from=start, candles=candles, today=today + ) return candles, "fresh_download" first_date, last_date = _date_bounds(cached) @@ -547,9 +706,18 @@ def ensure_daily_history( ) if not candles.empty: candles.to_parquet(path, index=False) + self._record_vendor_earliest( + symbol, security_id, requested_from=start, candles=candles, today=today + ) return candles, "fresh_download" - if first_date > start: + # A cache that starts late is either an interrupted prefetch (the vendor + # HAS the missing years, so refetch) or a stock that listed after the + # window opened (the vendor has nothing earlier, so refetching is waste + # forever). ``_vendor_earliest_for`` is what tells the two apart; without + # evidence it returns None and the original always-refetch rule applies. + vendor_earliest = self._vendor_earliest_for(symbol, security_id, start, today) + if not _cache_reaches_back_far_enough(first_date, start, vendor_earliest): # The cache may be current at the back but missing years at the # front, usually after an old interrupted prefetch. Refetch the # intended full window so long-lookback screeners see real history. @@ -560,10 +728,15 @@ def ensure_daily_history( from_date=start, to_date=today, ) + self._record_vendor_earliest( + symbol, security_id, requested_from=start, candles=candles, today=today + ) if not candles.empty: candles.to_parquet(path, index=False) return candles, "backfilled" return cached, "fresh" + # Falling through on purpose: a later listing still needs its daily + # top-up. Suppressing the backfill must not freeze the symbol's tail. if last_date >= today: return cached, "fresh" @@ -1195,17 +1368,18 @@ def cleanup_stale_cache_files( choose the age threshold, and only daily parquet files plus their sidecar markers are touched. - Two sidecar kinds travel with a parquet: `.checked` (the empty-increment - marker written here) and `.repaired` (the DATA-002 repair cooldown). Both - are meaningless without their parquet, so an orphan of either is removed - regardless of age. + Three sidecar kinds travel with a parquet: `.checked` (the empty-increment + marker written here), `.repaired` (the DATA-002 repair cooldown), and + `.firstbar` (the DATA-004 record of how far back the vendor's history + goes). All are meaningless without their parquet, so an orphan of any of + them is removed regardless of age. """ if not self.cache_dir.exists(): return 0 now = now or datetime.now() cutoff = now - timedelta(days=max(1, int(max_age_days))) targets: set[Path] = set() - sidecar_suffixes = (".checked", ".repaired") + sidecar_suffixes = (".checked", ".repaired", ".firstbar") for parquet in self.cache_dir.glob("*.parquet"): modified = datetime.fromtimestamp(parquet.stat().st_mtime) diff --git a/tests/test_daily_data_loader_vendor_earliest.py b/tests/test_daily_data_loader_vendor_earliest.py new file mode 100644 index 0000000..5d8968f --- /dev/null +++ b/tests/test_daily_data_loader_vendor_earliest.py @@ -0,0 +1,311 @@ +"""Tests for remembering how far back the vendor's history actually goes (DATA-004). + +Beginner note: +The cache-coverage test asks "does this file reach back to the start of the +requested window?". For a stock that listed *after* that start — DMART listed in +2017, well inside a ten-year window — the answer is permanently no, because DhanHQ +has nothing earlier to give. Before this change every prefetch and every scan +re-downloaded those symbols' whole history, wrote the same short frame back, and +did it again next time. 200 of 577 cached symbols were in that state. + +The fix records what the vendor actually served, following the ``.checked`` and +``.repaired`` sidecar precedent already in this codebase: once we have asked from +a given date and learned the earliest bar that exists, a cache reaching that bar is +as complete as it can ever be. + +The distinction these tests protect is the whole point: + +- an **interrupted prefetch** left a partial file and the vendor *does* have + earlier data → must still refetch; +- a **later listing** means the vendor has nothing earlier → refetching is waste. +""" + +from __future__ import annotations + +import json +from datetime import date, timedelta +from pathlib import Path +from typing import cast + +import pandas as pd + +from backend.daily_data_loader import DailyDataLoader +from backend.dhan_client import DhanDataClient + +TODAY = date(2026, 8, 24) +HISTORY_START = date(2016, 8, 24) # TODAY minus ten years +LISTED_ON = date(2017, 3, 21) # a DMART-style later listing +ROW = {"symbol": "DMART", "security_id": "1"} + + +def _month_series(first: date, last: date) -> pd.DatetimeIndex: + """Monthly bars that begin exactly on ``first`` and end on or before ``last``. + + ``pd.date_range(..., freq="MS")`` snaps to month starts, which would silently + move a mid-month listing date *and* leave the newest bar weeks short of the + requested end. Both endpoints are therefore inserted explicitly: these tests + turn on the first bar being exactly the listing date, and a snapped last bar + would make the frame fail the freshness half of the coverage test for reasons + that have nothing to do with what is being tested. + """ + months = pd.date_range(first, last, freq="MS") + return ( + pd.DatetimeIndex([pd.Timestamp(first), *months, pd.Timestamp(last)]) + .unique() + .sort_values() + ) + + +class ListedLateClient: + """A vendor that has no data before ``listed_on``, whatever you ask for. + + This is the real behaviour that made the coverage test unsatisfiable: the + request reaches back ten years, the response starts at the listing date. + """ + + def __init__(self, listed_on: date = LISTED_ON, *, through: date = TODAY) -> None: + self.listed_on = listed_on + self.through = through + # The loader passes real dates on every path these tests exercise. + self.calls: list[tuple[date, date]] = [] + + def fetch_daily_candles(self, *, from_date, to_date, **_kwargs) -> pd.DataFrame: + self.calls.append((from_date, to_date)) + # Monthly bars keep the fixture small; only the date bounds matter. The + # listing date itself is forced in because "MS" snaps to month starts, + # and the first bar being exactly the listing date is the whole point. + dates = _month_series(self.listed_on, self.through) + return pd.DataFrame( + { + "timestamp": dates, + "open": 100.0, + "high": 105.0, + "low": 99.0, + "close": 104.0, + "volume": 1_000.0, + } + ) + + +def _loader(tmp_path: Path, client: object) -> DailyDataLoader: + # The loader only calls fetch_daily_candles, so the duck-typed fake stands in. + return DailyDataLoader( + cast(DhanDataClient, client), cache_dir=tmp_path, request_delay_seconds=0.0 + ) + + +def _write_cache(loader: DailyDataLoader, first: date, last: date) -> Path: + path = loader.cache_path(ROW["symbol"], ROW["security_id"]) + pd.DataFrame( + { + "timestamp": _month_series(first, last), + "open": 100.0, + "high": 105.0, + "low": 99.0, + "close": 104.0, + "volume": 1_000.0, + } + ).to_parquet(path, index=False) + return path + + +def _marker(loader: DailyDataLoader) -> Path: + return loader.first_bar_path(ROW["symbol"], ROW["security_id"]) + + +# --------------------------------------------------------------------------- +# Learning the vendor's earliest bar +# --------------------------------------------------------------------------- + + +def test_a_full_download_records_the_vendors_earliest_bar(tmp_path: Path): + client = ListedLateClient() + loader = _loader(tmp_path, client) + + _frame, status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert status == "fresh_download" + payload = json.loads(_marker(loader).read_text(encoding="utf-8")) + # We asked from the ten-year start and learned the vendor begins at listing. + assert date.fromisoformat(payload["requested_from"]) == HISTORY_START + assert date.fromisoformat(payload["earliest_available"]) == LISTED_ON + + +def test_no_marker_is_written_when_the_vendor_covers_the_whole_window(tmp_path: Path): + """Nothing to remember when the request was fully satisfied.""" + client = ListedLateClient(listed_on=date(2016, 1, 1)) + loader = _loader(tmp_path, client) + + loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert not _marker(loader).exists() + + +def test_an_empty_response_records_nothing(tmp_path: Path): + """An empty answer is no evidence about how far back history goes.""" + + class EmptyClient: + def fetch_daily_candles(self, **_kwargs) -> pd.DataFrame: + return pd.DataFrame() + + loader = _loader(tmp_path, EmptyClient()) + + loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert not _marker(loader).exists() + + +# --------------------------------------------------------------------------- +# The bug: re-downloading a later listing forever +# --------------------------------------------------------------------------- + + +def test_a_later_listing_is_not_backfilled_again_on_the_next_pass(tmp_path: Path): + """The core defect. Second prefetch must not re-download the same history.""" + client = ListedLateClient() + loader = _loader(tmp_path, client) + + _frame, first_status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + calls_after_first = len(client.calls) + _frame, second_status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert first_status == "fresh_download" + assert second_status != "backfilled" + # The tail top-up may still ask; the ten-year backfill must not. + backfills = [c for c in client.calls[calls_after_first:] if c[0] == HISTORY_START] + assert backfills == [] + + +def test_a_later_listing_still_gets_its_daily_top_up(tmp_path: Path): + """Skipping the pointless backfill must not freeze the symbol's tail. + + The whole value of the cache is that it keeps advancing; suppressing the + backfill must only suppress the backfill. + """ + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON, last=TODAY - timedelta(days=40)) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY, + ) + + _frame, status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert status == "incremental" + # The request started after the cached tail, not at the ten-year start. + assert client.calls[-1][0] > TODAY - timedelta(days=41) + + +def test_a_later_listing_is_a_cache_hit_for_scans(tmp_path: Path): + """The scan path benefits too: no fetch, no rewrite.""" + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON, last=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY, + ) + + _frame, from_cache = loader.get_daily_history( + ROW, start_date=HISTORY_START, end_date=TODAY + ) + + assert from_cache is True + assert client.calls == [] + + +# --------------------------------------------------------------------------- +# The distinction that must not be lost +# --------------------------------------------------------------------------- + + +def test_a_genuinely_partial_cache_is_still_backfilled(tmp_path: Path): + """An interrupted prefetch, where the vendor DOES have earlier data. + + The marker says history begins in 2016, but the cache starts in 2020 — so + there are real bars missing and the backfill must run. + """ + client = ListedLateClient(listed_on=date(2016, 1, 1)) + loader = _loader(tmp_path, client) + _write_cache(loader, first=date(2020, 1, 1), last=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=date(2016, 1, 1), recorded_on=TODAY, + ) + + _frame, status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert status == "backfilled" + + +def test_a_shallower_probe_does_not_prove_anything_about_earlier_history(tmp_path: Path): + """A marker from a 5-year probe cannot answer a 10-year request. + + Learning that nothing exists before 2021 when you only asked from 2021 says + nothing about 2016, so the deeper request must still go to the vendor. + """ + client = ListedLateClient(listed_on=date(2021, 6, 1)) + loader = _loader(tmp_path, client) + _write_cache(loader, first=date(2021, 6, 1), last=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=date(2021, 1, 1), + earliest_available=date(2021, 6, 1), recorded_on=TODAY, + ) + + _frame, status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert status == "backfilled" + + +def test_a_stale_marker_is_re_probed(tmp_path: Path): + """Vendors do occasionally backfill history, so the belief expires.""" + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON, last=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=date(2026, 1, 1), + ) + + _frame, status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert status == "backfilled" + + +def test_an_unreadable_marker_fails_open_to_a_refetch(tmp_path: Path): + """A corrupt marker may cost a request; it must never hide missing history.""" + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON, last=TODAY) + _marker(loader).write_text("not json", encoding="utf-8") + + _frame, status = loader.ensure_daily_history(ROW, years_back=10, today=TODAY) + + assert status == "backfilled" + + +# --------------------------------------------------------------------------- +# Housekeeping +# --------------------------------------------------------------------------- + + +def test_orphan_first_bar_markers_are_cleaned_up(tmp_path: Path): + """`.firstbar` travels with its parquet, like `.checked` and `.repaired`.""" + loader = _loader(tmp_path, ListedLateClient()) + orphan = tmp_path / "GONE_9.firstbar" + orphan.write_text("{}", encoding="utf-8") + + removed = loader.cleanup_stale_cache_files(max_age_days=30) + + assert removed >= 1 + assert not orphan.exists() + + +def test_the_marker_suffix_is_distinct_from_the_other_sidecars(tmp_path: Path): + loader = _loader(tmp_path, ListedLateClient()) + path = _marker(loader) + + assert path.suffix == ".firstbar" + assert path.with_suffix(".checked") != path + assert path.with_suffix(".repaired") != path From 3e0104f69eeb7020f2d15c5138b7470413dc8f95 Mon Sep 17 00:00:00 2001 From: DoRmAmMu1997 Date: Fri, 4 Sep 2026 18:28:27 +0530 Subject: [PATCH 2/2] fix(DATA-004): harden vendor earliest evidence Co-authored-by: Codex --- backend/daily_data_loader.py | 156 ++++++++++-- .../components/data-acquisition.md | 32 +++ .../test_daily_data_loader_vendor_earliest.py | 222 +++++++++++++++++- 3 files changed, 387 insertions(+), 23 deletions(-) diff --git a/backend/daily_data_loader.py b/backend/daily_data_loader.py index 92a8545..d2bdc98 100644 --- a/backend/daily_data_loader.py +++ b/backend/daily_data_loader.py @@ -57,6 +57,79 @@ VENDOR_EARLIEST_RECHECK_DAYS = 30 +@dataclass(frozen=True) +class _VendorEarliestEvidence: + """Validated `.firstbar` evidence that bounds a vendor's history. + + Beginner note: + This small immutable object separates untrusted JSON on disk from the dates + the cache-coverage decision is allowed to trust. Once constructed, no caller + can accidentally change one date and leave the chronology inconsistent. + """ + + requested_from: date + earliest_available: date + recorded_on: date + + +def _read_vendor_earliest_evidence(path: Path) -> _VendorEarliestEvidence | None: + """Read one coherent `.firstbar` sidecar, or return no evidence. + + The sidecar is deliberately fail-open: malformed content cannot certify an + incomplete cache, so callers treat it as absent and refetch if necessary. + Its three dates must use canonical ``YYYY-MM-DD`` form and establish the + exact chronology ``requested_from < earliest_available <= recorded_on``. + Unknown JSON object fields are ignored to allow future metadata additions. + + Args: + path: Sidecar path next to the daily-cache parquet file. + + Returns: + Immutable evidence when all required fields are valid; otherwise ``None``. + + Beginner note: + ``date.fromisoformat`` also accepts compact and ISO-week strings. Checking + ``isoformat()`` after parsing prevents those alternate spellings from becoming + an undocumented on-disk format that future readers might interpret differently. + """ + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except (OSError, UnicodeDecodeError, json.JSONDecodeError): + return None + if not isinstance(payload, dict): + return None + + requested_from = payload.get("requested_from") + earliest_available = payload.get("earliest_available") + recorded_on = payload.get("recorded_on") + if ( + not isinstance(requested_from, str) + or not isinstance(earliest_available, str) + or not isinstance(recorded_on, str) + ): + return None + + try: + parsed_requested_from = date.fromisoformat(requested_from) + parsed_earliest_available = date.fromisoformat(earliest_available) + parsed_recorded_on = date.fromisoformat(recorded_on) + except ValueError: + return None + if ( + parsed_requested_from.isoformat() != requested_from + or parsed_earliest_available.isoformat() != earliest_available + or parsed_recorded_on.isoformat() != recorded_on + ): + return None + if not parsed_requested_from < parsed_earliest_available <= parsed_recorded_on: + return None + return _VendorEarliestEvidence( + requested_from=parsed_requested_from, + earliest_available=parsed_earliest_available, + recorded_on=parsed_recorded_on, + ) + + def history_start_date( years_back: int = DEFAULT_HISTORY_YEARS_BACK, today: date | None = None ) -> date: @@ -235,15 +308,16 @@ def _cache_reaches_back_far_enough( ``vendor_earliest`` is the earliest bar the vendor actually served for a probe that reached at least as far back as this request (see - ``DailyDataLoader._vendor_earliest_for``). When the cache already starts at or - before that bar, it is as complete as it can ever be, so asking again is pure - waste. ``None`` means we have no such evidence and the strict rule applies — - which keeps an interrupted prefetch's partial file being refetched, since - there the vendor genuinely does have the missing years. + ``DailyDataLoader._vendor_earliest_for``). It must match the cache's literal + first date: accepting an earlier or later cache start would let contradictory + evidence certify an incomplete or corrupted cache. ``None`` means we have no + such evidence and the strict rule applies — which keeps an interrupted + prefetch's partial file being refetched, since there the vendor genuinely has + the missing years. """ if first_date <= requested_start: return True - return vendor_earliest is not None and first_date <= vendor_earliest + return vendor_earliest is not None and first_date == vendor_earliest def _date_bounds(candles: pd.DataFrame) -> tuple[date | None, date | None]: @@ -483,22 +557,29 @@ def _vendor_earliest_for( Any read or parse problem returns ``None``, so a corrupt marker can only ever cost an extra request — never hide history that really is missing. + + Args: + symbol: Instrument symbol used to locate the sidecar. + security_id: Vendor identifier paired with the symbol in cache paths. + requested_start: Earliest date the current caller needs covered. + + Returns: + The qualifying earliest available date, or ``None`` for no authority. + + Beginner note: + A marker that is too new in the future is no more trustworthy than a stale + one: both have a clock relationship that cannot describe a completed + vendor request, so the cache takes the safe (refetch) path. """ - path = self.first_bar_path(symbol, security_id) - if not path.exists(): - return None - try: - payload = json.loads(path.read_text(encoding="utf-8")) - requested_from = _coerce_date(str(payload["requested_from"])) - earliest_available = _coerce_date(str(payload["earliest_available"])) - recorded_on = _coerce_date(str(payload["recorded_on"])) - except (OSError, ValueError, KeyError, TypeError, json.JSONDecodeError): + evidence = _read_vendor_earliest_evidence(self.first_bar_path(symbol, security_id)) + if evidence is None: return None - if (self.today_func() - recorded_on).days >= VENDOR_EARLIEST_RECHECK_DAYS: + age_days = (self.today_func() - evidence.recorded_on).days + if age_days < 0 or age_days >= VENDOR_EARLIEST_RECHECK_DAYS: return None - if requested_from > requested_start: + if evidence.requested_from > requested_start: return None - return earliest_available + return evidence.earliest_available def _write_vendor_earliest( self, @@ -535,21 +616,52 @@ def _record_vendor_earliest( ) -> None: """Record the vendor's earliest bar when it fell short of what we asked for. - Called after any full-window download. A response that *does* reach the - requested start teaches us nothing worth storing, and an empty response is - no evidence at all, so both are skipped. + Called after any full-window download. Empty or invalid frames are + inconclusive and leave an existing marker untouched. A shallower probe + cannot replace or renew fresh deeper evidence. Expired or future-dated + evidence is not authoritative and may be replaced. An equally + deep/deeper response that reaches the requested start invalidates its + old marker; otherwise a later first bar becomes new evidence. ``recorded_on`` is stamped from the injected wall clock rather than from the caller's requested window, so the marker ages in real time whatever range was asked for. + + Args: + symbol: Instrument symbol used to locate the sidecar. + security_id: Vendor identifier paired with the symbol in cache paths. + requested_from: Inclusive beginning of the completed vendor probe. + candles: Raw non-empty response used to learn its first bar. + + Beginner note: + The depth comparison is between what the vendor was asked, not what it + returned. A late-listed stock can return the same first bar for many + windows, but only the request that began furthest back proves the stronger + "nothing exists earlier" statement. """ + start = _coerce_date(requested_from) if candles.empty: return first_date, _last_date = _date_bounds(candles) if first_date is None: return - start = _coerce_date(requested_from) + + path = self.first_bar_path(symbol, security_id) + existing = _read_vendor_earliest_evidence(path) + # A probe beginning later asks less of the vendor. Preserve evidence + # collected by a deeper request only while its wall-clock TTL is valid; + # expired/future-dated evidence is not authoritative and must not block + # this fresh answer from replacing it. + if existing is not None and existing.requested_from < start: + age_days = (self.today_func() - existing.recorded_on).days + if 0 <= age_days < VENDOR_EARLIEST_RECHECK_DAYS: + return if first_date <= start: + try: + path.unlink(missing_ok=True) + except OSError: + # The marker is optional, so a locked sidecar must not fail a fetch. + logger.warning("Could not remove obsolete daily-cache first-bar marker for %s", symbol) return self._write_vendor_earliest( symbol, diff --git a/docs/architecture/components/data-acquisition.md b/docs/architecture/components/data-acquisition.md index 20e39a1..bdeecfd 100644 --- a/docs/architecture/components/data-acquisition.md +++ b/docs/architecture/components/data-acquisition.md @@ -67,6 +67,38 @@ flowchart TD | `history_start_date(years_back, today)` | Leap-safe "subtract whole years" (Feb 29 → Feb 28). | | `safe_file_stem(value)` | Path-traversal-safe filename fragment. | +### DATA-004 `.firstbar` earliest-history evidence + +When a vendor request begins before a stock listed, the returned frame begins at +the stock's earliest available candle. `DailyDataLoader` stores that answer next +to the parquet as `_.firstbar`, a JSON object with canonical +`requested_from`, `earliest_available`, and `recorded_on` dates. The public cache +contract remains strict: a request is front-complete only when the parquet first +date literally reaches `requested_start`, or a fresh qualifying `.firstbar` +exists **and its `earliest_available` exactly equals the parquet first date**. +The exact binding prevents contradictory cache/sidecar dates from certifying +history that cannot be proved complete. + +The marker is internal, optional evidence—not a user input or a database record. +Its JSON reader accepts only a JSON object with string `YYYY-MM-DD` fields and +the chronology `requested_from < earliest_available <= recorded_on`; extra fields +are ignored for forwards compatibility. Its 30-day TTL is measured against the +injected wall clock, not the requested data window. Future, stale (age 30 days or +more), malformed, noncanonical, or shallower-than-request evidence is ignored and +therefore causes a safe refetch. A shallow probe preserves a deeper marker only +while that marker is fresh; expired or future-dated evidence is replaced by the +new probe, without renewing a fresh marker's timestamp. An equally deep/deeper +non-empty response that still starts late replaces the marker; one that reaches +the requested start removes the now-obsolete marker best-effort. Empty/invalid +frames leave prior evidence unchanged. + +`.firstbar` shares the cache lifecycle: it travels with its parquet and +`cleanup_stale_cache_files()` removes it when orphaned or when the associated +cache ages out. `tests/test_daily_data_loader_vendor_earliest.py` covers marker +creation, strict parsing/chronology, exact cache binding, wall-clock TTL, +shallower/equally-deep update rules, cleanup, and the request-bounded late-listing +vendor fixture. + ## 4. Key design decisions & trade-offs | Decision | Rationale | Alternative rejected | diff --git a/tests/test_daily_data_loader_vendor_earliest.py b/tests/test_daily_data_loader_vendor_earliest.py index 9ec10f1..f296e17 100644 --- a/tests/test_daily_data_loader_vendor_earliest.py +++ b/tests/test_daily_data_loader_vendor_earliest.py @@ -28,6 +28,7 @@ from typing import cast import pandas as pd +import pytest from backend.daily_data_loader import DailyDataLoader from backend.dhan_client import DhanDataClient @@ -74,7 +75,10 @@ def fetch_daily_candles(self, *, from_date, to_date, **_kwargs) -> pd.DataFrame: # Monthly bars keep the fixture small; only the date bounds matter. The # listing date itself is forced in because "MS" snaps to month starts, # and the first bar being exactly the listing date is the whole point. - dates = _month_series(self.listed_on, self.through) + first_returned_date = max(self.listed_on, from_date) + if first_returned_date > self.through: + return pd.DataFrame() + dates = _month_series(first_returned_date, self.through) return pd.DataFrame( { "timestamp": dates, @@ -379,3 +383,219 @@ def test_a_recent_marker_is_honoured_against_the_clock(tmp_path: Path): assert loader._vendor_earliest_for( ROW["symbol"], ROW["security_id"], HISTORY_START ) == LISTED_ON + + +# --------------------------------------------------------------------------- +# Marker hardening: evidence is useful only when it is self-consistent +# --------------------------------------------------------------------------- + + +def test_a_future_dated_marker_fails_open_to_a_backfill(tmp_path: Path): + """Future evidence must not certify a cache before that day has arrived. + + Beginner note: + The marker is only an optimisation. Treating a future stamp as fresh would + let a clock error or manually edited sidecar hide missing history, whereas + rejecting it merely asks the vendor again. + """ + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON, last=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY + timedelta(days=1), + ) + + _frame, from_cache = loader.get_daily_history( + ROW, start_date=HISTORY_START, end_date=TODAY + ) + + assert from_cache is False + assert client.calls == [(HISTORY_START, TODAY)] + + +def test_invalid_utf8_marker_fails_open_to_a_backfill(tmp_path: Path): + """Undecodable marker bytes must trigger a safe refetch, not escape parsing.""" + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON, last=TODAY) + _marker(loader).write_bytes(b"\xff") + + _frame, from_cache = loader.get_daily_history( + ROW, start_date=HISTORY_START, end_date=TODAY + ) + + assert from_cache is False + assert client.calls == [(HISTORY_START, TODAY)] + + +def test_a_marker_must_match_the_cached_first_bar_to_avoid_a_backfill(tmp_path: Path): + """Contradictory cache and marker starts must trigger a refetch.""" + client = ListedLateClient() + loader = _loader(tmp_path, client) + _write_cache(loader, first=LISTED_ON - timedelta(days=1), last=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY, + ) + + _frame, from_cache = loader.get_daily_history( + ROW, start_date=HISTORY_START, end_date=TODAY + ) + + assert from_cache is False + assert client.calls == [(HISTORY_START, TODAY)] + + +def test_noncanonical_non_string_or_impossible_marker_fields_are_rejected(tmp_path: Path): + """The parser accepts only a coherent object with canonical ISO dates.""" + loader = _loader(tmp_path, ListedLateClient()) + invalid_payloads = [ + { + "requested_from": 20160824, + "earliest_available": LISTED_ON.isoformat(), + "recorded_on": TODAY.isoformat(), + }, + { + "requested_from": "20160824", + "earliest_available": LISTED_ON.isoformat(), + "recorded_on": TODAY.isoformat(), + }, + { + "requested_from": HISTORY_START.isoformat(), + "earliest_available": HISTORY_START.isoformat(), + "recorded_on": TODAY.isoformat(), + }, + { + "requested_from": HISTORY_START.isoformat(), + "earliest_available": (TODAY + timedelta(days=1)).isoformat(), + "recorded_on": TODAY.isoformat(), + }, + ] + + for payload in invalid_payloads: + _marker(loader).write_text(json.dumps(payload), encoding="utf-8") + + assert loader._vendor_earliest_for( + ROW["symbol"], ROW["security_id"], HISTORY_START + ) is None + + +def test_a_fresh_deeper_marker_survives_a_shallower_qualifying_probe(tmp_path: Path): + """A shallower response cannot replace or renew fresh stronger evidence.""" + loader = _loader(tmp_path, ListedLateClient()) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY - timedelta(days=2), + ) + original = _marker(loader).read_text(encoding="utf-8") + shallow_start = HISTORY_START + timedelta(days=100) + shallow_candles = pd.DataFrame( + {"timestamp": _month_series(date(2020, 1, 1), TODAY)} + ) + + loader._record_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=shallow_start, + candles=shallow_candles, + ) + + assert _marker(loader).read_text(encoding="utf-8") == original + + +@pytest.mark.parametrize( + "recorded_on", + [TODAY - timedelta(days=30), TODAY + timedelta(days=1)], + ids=["expired", "future-dated"], +) +def test_expired_or_future_deeper_marker_does_not_block_a_shallower_probe( + tmp_path: Path, recorded_on: date +): + """Non-fresh deeper evidence must not suppress a new shallower answer. + + Beginner note: + A deeper marker is stronger only while its 30-day wall-clock TTL is valid. + Once it is expired or dated in the future, keeping it would let old or + clock-skewed evidence block a new probe forever. The new probe therefore + replaces it, just as it would when no marker existed. + """ + loader = _loader(tmp_path, ListedLateClient()) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=recorded_on, + ) + shallow_start = HISTORY_START + timedelta(days=100) + shallow_candles = pd.DataFrame( + {"timestamp": _month_series(date(2020, 1, 1), TODAY)} + ) + + loader._record_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=shallow_start, + candles=shallow_candles, + ) + + payload = json.loads(_marker(loader).read_text(encoding="utf-8")) + assert payload["requested_from"] == shallow_start.isoformat() + assert payload["earliest_available"] == date(2020, 1, 1).isoformat() + assert payload["recorded_on"] == TODAY.isoformat() + + +def test_a_full_equally_deep_probe_removes_obsolete_marker(tmp_path: Path): + """A response reaching the requested start makes its old marker obsolete.""" + loader = _loader(tmp_path, ListedLateClient()) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY, + ) + complete_candles = pd.DataFrame( + {"timestamp": _month_series(HISTORY_START, TODAY)} + ) + + loader._record_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + candles=complete_candles, + ) + + assert not _marker(loader).exists() + + +def test_marker_age_29_days_is_fresh_but_age_30_days_is_rejected(tmp_path: Path): + """The 30-day wall-clock TTL is inclusive at its expiration boundary.""" + loader = _loader(tmp_path, ListedLateClient()) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY - timedelta(days=29), + ) + + assert loader._vendor_earliest_for( + ROW["symbol"], ROW["security_id"], HISTORY_START + ) == LISTED_ON + + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY - timedelta(days=30), + ) + + assert loader._vendor_earliest_for( + ROW["symbol"], ROW["security_id"], HISTORY_START + ) is None + + +def test_listed_late_client_respects_the_requested_start_and_end_dates() -> None: + """The fixture must model a bounded vendor response, not a full cache read.""" + client = ListedLateClient() + requested_start = TODAY - timedelta(days=10) + frame = client.fetch_daily_candles(from_date=requested_start, to_date=TODAY) + + assert frame["timestamp"].min().date() == requested_start + assert frame["timestamp"].max().date() == TODAY + + +def test_listed_late_client_returns_no_rows_after_its_configured_through_date() -> None: + """A request beginning after the fixture's vendor horizon is empty evidence.""" + client = ListedLateClient() + + frame = client.fetch_daily_candles( + from_date=TODAY + timedelta(days=1), to_date=TODAY + timedelta(days=2) + ) + + assert frame.empty