diff --git a/backend/daily_data_loader.py b/backend/daily_data_loader.py index 728d32f..d2bdc98 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 @@ -48,6 +49,86 @@ # ``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 + + +@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 @@ -155,8 +236,9 @@ def _cache_covers_range( last_date: date | None, requested_start: date, requested_end: date, - checked_through: date | None = None, *, + vendor_earliest: date | None = None, + checked_through: date | None = None, allow_unpublished_tail: bool = False, ) -> bool: """Return whether a cached range has the authority to answer a request. @@ -176,6 +258,7 @@ def _cache_covers_range( last_date: Newest valid candle date in the candidate cache. requested_start: Inclusive requested start date. requested_end: Inclusive requested end date. + vendor_earliest: Earliest bar a qualifying vendor probe returned. checked_through: Optional prefetch sidecar date recording an empty tail. allow_unpublished_tail: Scanner-only authority to use current-session and sidecar-marker tail relaxation. @@ -188,7 +271,9 @@ def _cache_covers_range( """ if first_date is None or last_date is None: return False - if first_date > requested_start: + # Front and back are judged by separate evidence: how far the vendor's history + # goes (DATA-004) and which recent days it has actually published (DATA-003). + if not _cache_reaches_back_far_enough(first_date, requested_start, vendor_earliest): return False if _only_unpublished_days_missing( last_date, @@ -208,6 +293,33 @@ def _cache_covers_range( return checked_through is not None and checked_through >= requested_end +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``). 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 + + def _date_bounds(candles: pd.DataFrame) -> tuple[date | None, date | None]: """Return the first/last valid candle dates in a cached frame. @@ -322,6 +434,7 @@ def __init__( max_consecutive_failures: int | None = None, sleep_func: Callable[[float], None] = time.sleep, fetch_workers: int | None = None, + today_func: Callable[[], date] = date.today, ): # The Dhan client is optional so cache-only callers (the legacy-file # cleanup step, the chart UI's `read_cached_history`) can build a loader @@ -347,6 +460,12 @@ def __init__( ) self.max_consecutive_failures = max(0, int(max_consecutive_failures or 0)) self.sleep_func = sleep_func + # Wall clock, injected like sleep_func so tests can pin it. Deliberately + # separate from the ``today`` argument callers pass to describe a DATA + # window: "the date I am asking about" and "the date it is now" are + # different questions, and conflating them let a marker's 30-day expiry be + # judged against a historical request boundary, so it never expired. + self.today_func = today_func # PERF-001: 1 (the default) keeps the long-standing sequential path # byte-identical. Values above 1 fetch with a thread pool while the # shared pacer holds the global inter-request delay. @@ -408,6 +527,150 @@ 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 + ) -> 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`` of **now** + (vendors do backfill occasionally, so the belief expires). Age is + measured against the injected wall clock, never against the requested + window: judging it by a request boundary meant a repeated historical + request always computed an age of zero and the marker never expired. + - 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. + + 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. + """ + evidence = _read_vendor_earliest_evidence(self.first_bar_path(symbol, security_id)) + if evidence is None: + return None + age_days = (self.today_func() - evidence.recorded_on).days + if age_days < 0 or age_days >= VENDOR_EARLIEST_RECHECK_DAYS: + return None + if evidence.requested_from > requested_start: + return None + return evidence.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, + ) -> None: + """Record the vendor's earliest bar when it fell short of what we asked for. + + 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 + + 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, + security_id, + requested_from=start, + earliest_available=first_date, + recorded_on=self.today_func(), + ) + 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. @@ -503,6 +766,9 @@ def get_daily_history( first_date, last_date = _date_bounds(cached) requested_start = _coerce_date(start_date) requested_end = _coerce_date(end_date) + vendor_earliest = self._vendor_earliest_for( + symbol, security_id, requested_start + ) # The prefetch's marker is the only evidence that distinguishes a # market holiday (a weekday with no bar to fetch) from a weekday whose # bar we simply have not collected yet. @@ -512,7 +778,8 @@ def get_daily_history( last_date, requested_start, requested_end, - checked_through, + vendor_earliest=vendor_earliest, + checked_through=checked_through, allow_unpublished_tail=allow_unpublished_tail, ): if cached is None: @@ -526,7 +793,8 @@ def get_daily_history( actual_last, requested_start, requested_end, - checked_through, + vendor_earliest=vendor_earliest, + checked_through=checked_through, allow_unpublished_tail=allow_unpublished_tail, ): return self._slice_to_range(cached, start_date, end_date), True @@ -543,6 +811,9 @@ 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 + ) return self._slice_to_range(candles, start_date, end_date), False def fetch_window( @@ -624,6 +895,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 + ) return candles, "fresh_download" cached = pd.read_parquet(path) @@ -639,6 +913,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 + ) return candles, "fresh_download" first_date, last_date = _date_bounds(cached) @@ -655,9 +932,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 + ) 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) + 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. @@ -668,10 +954,15 @@ def ensure_daily_history( from_date=start, to_date=today, ) + self._record_vendor_earliest( + symbol, security_id, requested_from=start, candles=candles + ) 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" @@ -1318,17 +1609,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/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 new file mode 100644 index 0000000..f296e17 --- /dev/null +++ b/tests/test_daily_data_loader_vendor_earliest.py @@ -0,0 +1,601 @@ +"""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 +import pytest + +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. + 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, + "open": 100.0, + "high": 105.0, + "low": 99.0, + "close": 104.0, + "volume": 1_000.0, + } + ) + + +def _loader( + tmp_path: Path, client: object, *, now: date = TODAY +) -> DailyDataLoader: + """Build a loader with a pinned wall clock. + + ``now`` is the injected clock, which is what the marker's expiry is measured + against. It is deliberately separate from the ``today`` argument these tests + pass to ``ensure_daily_history``: that one describes the *data* window. + """ + # 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, + today_func=lambda: now, + ) + + +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 + + +# --------------------------------------------------------------------------- +# The marker ages in real time, not against the requested window +# --------------------------------------------------------------------------- + + +def test_a_marker_from_a_historical_request_still_expires(tmp_path: Path): + """Codex review, PR #114. + + ``get_daily_history`` used to stamp ``recorded_on`` with the request's own + ``end_date``. Repeating a historical request then computed an age of zero every + time, so the 30-day expiry never fired and a vendor backfill could stay hidden + behind a partial cache indefinitely. + """ + historical_end = date(2020, 6, 30) + client = ListedLateClient(listed_on=LISTED_ON, through=historical_end) + # The clock is far ahead of the window being requested. + loader = _loader(tmp_path, client, now=TODAY) + + loader.get_daily_history(ROW, start_date=HISTORY_START, end_date=historical_end) + + payload = json.loads(_marker(loader).read_text(encoding="utf-8")) + # Stamped from the clock, not from the 2020 request boundary. + assert date.fromisoformat(payload["recorded_on"]) == TODAY + # And it is therefore already long expired relative to when it was "recorded". + stale_loader = _loader(tmp_path, client, now=TODAY + timedelta(days=31)) + assert stale_loader._vendor_earliest_for( + ROW["symbol"], ROW["security_id"], HISTORY_START + ) is None + + +def test_a_future_dated_request_does_not_extend_a_markers_life(tmp_path: Path): + """A request reaching into the future must not keep a stale belief alive.""" + client = ListedLateClient() + loader = _loader(tmp_path, client, now=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY - timedelta(days=45), + ) + + # Asking about a window that ends next year must not make a 45-day-old + # marker look fresh. + assert loader._vendor_earliest_for( + ROW["symbol"], ROW["security_id"], HISTORY_START + ) is None + + +def test_a_recent_marker_is_honoured_against_the_clock(tmp_path: Path): + """The positive case, so the expiry test above cannot pass vacuously.""" + loader = _loader(tmp_path, ListedLateClient(), now=TODAY) + loader._write_vendor_earliest( + ROW["symbol"], ROW["security_id"], requested_from=HISTORY_START, + earliest_available=LISTED_ON, recorded_on=TODAY - timedelta(days=5), + ) + + 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