diff --git a/data/cache/coverage.py b/data/cache/coverage.py index 29d6933..dd226be 100644 --- a/data/cache/coverage.py +++ b/data/cache/coverage.py @@ -52,23 +52,58 @@ class CoverageLedger: - """Append-only parquet ledger of fetched (endpoint, key, range) coverage.""" + """Append-only parquet ledger of fetched (endpoint, key, range) coverage. + + D4: a per-instance frame cache + lookup memos avoid re-reading and re-filtering + the full parquet on every repeated lookup. The cache is invalidated on the + file's mtime change (so an out-of-instance write is not served stale) and is + refreshed in place on writes through this instance. The public columns, parquet + path, and coverage semantics are unchanged. + """ def __init__(self, root: str, schema_version: str = "v1") -> None: self._root = Path(root) self._schema_version = schema_version + # in-process cache of the full ledger frame + the mtime it was read at + # (None mtime == file absent). Lookup results are memoized per (endpoint, + # key) and cleared whenever the frame cache is invalidated/refreshed. + self._frame: pd.DataFrame | None = None + self._frame_mtime: int | None = None + self._intervals_memo: dict[tuple[str, str], list[Interval]] = {} + self._snapshot_memo: dict[tuple[str, str], pd.Timestamp | None] = {} @property def path(self) -> Path: return self._root / "manifest" / "coverage.parquet" - # -- read --------------------------------------------------------------- # - def read(self) -> pd.DataFrame: - """Return the full ledger (empty, correctly-typed frame if absent).""" + # -- internal load path (mtime-invalidated cache) ----------------------- # + def _read_frame(self) -> pd.DataFrame: + """Read the ledger parquet from disk (the cache-miss path; spy target).""" if not self.path.exists(): return pd.DataFrame(columns=LEDGER_COLUMNS) return pd.read_parquet(self.path) + def _current_mtime(self) -> int | None: + try: + return self.path.stat().st_mtime_ns + except FileNotFoundError: + return None + + def _load(self) -> pd.DataFrame: + """Return the cached full ledger, reloading only when the file changed.""" + mtime = self._current_mtime() + if self._frame is None or mtime != self._frame_mtime: + self._frame = self._read_frame() + self._frame_mtime = mtime + self._intervals_memo.clear() + self._snapshot_memo.clear() + return self._frame + + # -- read --------------------------------------------------------------- # + def read(self) -> pd.DataFrame: + """Return the full ledger (a COPY, so callers cannot corrupt the cache).""" + return self._load().copy() + def covered_intervals(self, endpoint: str, key: str) -> list[Interval]: """Closed date intervals already covered for ``(endpoint, key)``. @@ -76,8 +111,12 @@ def covered_intervals(self, endpoint: str, key: str) -> list[Interval]: Snapshot rows (NaT range) are ignored here — this is the dense-endpoint planner's view (P4-1 only uses dense endpoints). """ - ledger = self.read() + ledger = self._load() # first, so an external change invalidates the memo + memo_key = (endpoint, key) + if memo_key in self._intervals_memo: + return list(self._intervals_memo[memo_key]) if ledger.empty: + self._intervals_memo[memo_key] = [] return [] mask = ( (ledger["endpoint"] == endpoint) @@ -87,10 +126,12 @@ def covered_intervals(self, endpoint: str, key: str) -> list[Interval]: & ledger["end_date"].notna() ) sub = ledger[mask] - return [ + result = [ (pd.Timestamp(s), pd.Timestamp(e)) for s, e in zip(sub["start_date"], sub["end_date"]) ] + self._intervals_memo[memo_key] = result + return list(result) def snapshot_fetched_at(self, endpoint: str, key: str) -> pd.Timestamp | None: """Latest successful fetch time for a snapshot/dimension ``(endpoint, key)``. @@ -100,20 +141,42 @@ def snapshot_fetched_at(self, endpoint: str, key: str) -> pd.Timestamp | None: ``empty`` ``fetched_at``. ``None`` means never fetched (must fetch). The caller compares against ``refresh_dimension_days`` to decide staleness. """ - ledger = self.read() - if ledger.empty: - return None - mask = ( - (ledger["endpoint"] == endpoint) - & (ledger["key"] == str(key)) - & (ledger["status"].isin(_COVERING_STATUSES)) - ) - sub = ledger[mask] - if sub.empty: - return None - return pd.Timestamp(sub["fetched_at"].max()) + ledger = self._load() # first, so an external change invalidates the memo + memo_key = (endpoint, str(key)) + if memo_key in self._snapshot_memo: + return self._snapshot_memo[memo_key] + result: pd.Timestamp | None = None + if not ledger.empty: + mask = ( + (ledger["endpoint"] == endpoint) + & (ledger["key"] == str(key)) + & (ledger["status"].isin(_COVERING_STATUSES)) + ) + sub = ledger[mask] + if not sub.empty: + result = pd.Timestamp(sub["fetched_at"].max()) + self._snapshot_memo[memo_key] = result + return result # -- write -------------------------------------------------------------- # + def _normalize_row(self, row: dict) -> dict: + """Normalize one input row to the ledger dtypes/order (no secret fields).""" + start_date = row.get("start_date") + end_date = row.get("end_date") + return { + "endpoint": row["endpoint"], + "key_type": row["key_type"], + "key": str(row["key"]), + "start_date": pd.Timestamp(start_date) if start_date is not None else pd.NaT, + "end_date": pd.Timestamp(end_date) if end_date is not None else pd.NaT, + "fields_hash": str(row["fields_hash"]), + "fetched_at": pd.Timestamp(row["fetched_at"]), + "row_count": int(row["row_count"]), + "status": str(row["status"]), + "schema_version": self._schema_version, + "source_version": row.get("source_version"), + } + def record( self, *, @@ -128,26 +191,41 @@ def record( fetched_at: pd.Timestamp, source_version: str | None = None, ) -> None: - """Append one coverage row (atomic). No secret is ever written here.""" - row = { - "endpoint": endpoint, - "key_type": key_type, - "key": str(key), - "start_date": pd.Timestamp(start_date) if start_date is not None else pd.NaT, - "end_date": pd.Timestamp(end_date) if end_date is not None else pd.NaT, - "fields_hash": str(fields_hash), - "fetched_at": pd.Timestamp(fetched_at), - "row_count": int(row_count), - "status": str(status), - "schema_version": self._schema_version, - "source_version": source_version, - } - existing = self.read() - new_row = pd.DataFrame([row]) + """Append one coverage row. Delegates to :meth:`record_many` so the + single-row path stays identical to the batch path.""" + self.record_many([ + { + "endpoint": endpoint, + "key_type": key_type, + "key": key, + "start_date": start_date, + "end_date": end_date, + "fields_hash": fields_hash, + "row_count": row_count, + "status": status, + "fetched_at": fetched_at, + "source_version": source_version, + } + ]) + + def record_many(self, rows: list[dict]) -> None: + """Append a batch of coverage rows with one normalize + one parquet write. + + Empty input is a no-op. Each row is normalized exactly like a single + ``record``; the written parquet is reindexed to ``LEDGER_COLUMNS``. No + token / secret is ever written here. The in-process cache is refreshed to + the just-written frame. + """ + rows = list(rows) + if not rows: + return + normalized = [self._normalize_row(r) for r in rows] + existing = self._load() + new_rows = pd.DataFrame(normalized) # Avoid concat-with-empty (it warns and can shift dtypes): the first - # record IS the new row; later records append to a non-empty ledger. - combined = new_row if existing.empty else pd.concat( - [existing, new_row], ignore_index=True + # records ARE the new rows; later records append to a non-empty ledger. + combined = new_rows if existing.empty else pd.concat( + [existing, new_rows], ignore_index=True ) # enforce column order / presence combined = combined.reindex(columns=LEDGER_COLUMNS) @@ -155,3 +233,8 @@ def record( tmp = self.path.with_suffix(".parquet.tmp") combined.to_parquet(tmp, engine="pyarrow", index=False) os.replace(tmp, self.path) + # refresh the in-process cache to the just-written frame (no re-read) + self._frame = combined + self._frame_mtime = self._current_mtime() + self._intervals_memo.clear() + self._snapshot_memo.clear() diff --git a/data/cache/intraday_coverage.py b/data/cache/intraday_coverage.py index 63436d4..bd80204 100644 --- a/data/cache/intraday_coverage.py +++ b/data/cache/intraday_coverage.py @@ -55,23 +55,53 @@ class IntradayCoverageLedger: - """Append-only parquet ledger of fetched (endpoint, symbol, time-range) coverage.""" + """Append-only parquet ledger of fetched (endpoint, symbol, time-range) coverage. + + D4: mirrors :class:`data.cache.coverage.CoverageLedger` — a per-instance frame + cache + lookup memo avoid re-reading and re-filtering the full parquet on every + repeated ``covered_day_intervals`` lookup, invalidated on the file's mtime + change and refreshed on writes through this instance. Public columns, parquet + path, and coverage semantics are unchanged. + """ def __init__(self, root: str, schema_version: str = "v1") -> None: self._root = Path(root) self._schema_version = schema_version + self._frame: pd.DataFrame | None = None + self._frame_mtime: int | None = None + self._day_memo: dict[tuple[str, str, str], list[Interval]] = {} @property def path(self) -> Path: return self._root / "manifest" / "coverage_intraday.parquet" - # -- read --------------------------------------------------------------- # - def read(self) -> pd.DataFrame: - """Return the full ledger (empty, correctly-typed frame if absent).""" + # -- internal load path (mtime-invalidated cache) ----------------------- # + def _read_frame(self) -> pd.DataFrame: + """Read the ledger parquet from disk (the cache-miss path; spy target).""" if not self.path.exists(): return pd.DataFrame(columns=INTRADAY_LEDGER_COLUMNS) return pd.read_parquet(self.path) + def _current_mtime(self) -> int | None: + try: + return self.path.stat().st_mtime_ns + except FileNotFoundError: + return None + + def _load(self) -> pd.DataFrame: + """Return the cached full ledger, reloading only when the file changed.""" + mtime = self._current_mtime() + if self._frame is None or mtime != self._frame_mtime: + self._frame = self._read_frame() + self._frame_mtime = mtime + self._day_memo.clear() + return self._frame + + # -- read --------------------------------------------------------------- # + def read(self) -> pd.DataFrame: + """Return the full ledger (a COPY, so callers cannot corrupt the cache).""" + return self._load().copy() + def covered_day_intervals( self, endpoint: str, key: str, raw_freq: str ) -> list[Interval]: @@ -82,8 +112,12 @@ def covered_day_intervals( each covered span is reported day-normalized for the day-interval algebra in :mod:`data.cache.intervals`. Only ``ok``/``empty`` rows count. """ - ledger = self.read() + ledger = self._load() # first, so an external change invalidates the memo + memo_key = (endpoint, str(key), str(raw_freq)) + if memo_key in self._day_memo: + return list(self._day_memo[memo_key]) if ledger.empty: + self._day_memo[memo_key] = [] return [] mask = ( (ledger["endpoint"] == endpoint) @@ -94,12 +128,32 @@ def covered_day_intervals( & ledger["end_time"].notna() ) sub = ledger[mask] - return [ + result = [ (pd.Timestamp(s).normalize(), pd.Timestamp(e).normalize()) for s, e in zip(sub["start_time"], sub["end_time"]) ] + self._day_memo[memo_key] = result + return list(result) # -- write -------------------------------------------------------------- # + def _normalize_row(self, row: dict) -> dict: + """Normalize one input row to the ledger dtypes/order (no secret fields).""" + start_time = row.get("start_time") + end_time = row.get("end_time") + return { + "endpoint": row["endpoint"], + "key_type": row["key_type"], + "key": str(row["key"]), + "raw_freq": str(row["raw_freq"]), + "start_time": pd.Timestamp(start_time) if start_time is not None else pd.NaT, + "end_time": pd.Timestamp(end_time) if end_time is not None else pd.NaT, + "fields_hash": str(row["fields_hash"]), + "fetched_at": pd.Timestamp(row["fetched_at"]), + "row_count": int(row["row_count"]), + "status": str(row["status"]), + "schema_version": self._schema_version, + } + def record( self, *, @@ -114,27 +168,45 @@ def record( status: str, fetched_at: pd.Timestamp, ) -> None: - """Append one coverage row (atomic). No secret is ever written here.""" - row = { - "endpoint": endpoint, - "key_type": key_type, - "key": str(key), - "raw_freq": str(raw_freq), - "start_time": pd.Timestamp(start_time) if start_time is not None else pd.NaT, - "end_time": pd.Timestamp(end_time) if end_time is not None else pd.NaT, - "fields_hash": str(fields_hash), - "fetched_at": pd.Timestamp(fetched_at), - "row_count": int(row_count), - "status": str(status), - "schema_version": self._schema_version, - } - existing = self.read() - new_row = pd.DataFrame([row]) - combined = new_row if existing.empty else pd.concat( - [existing, new_row], ignore_index=True + """Append one coverage row. Delegates to :meth:`record_many` so the + single-row path stays identical to the batch path.""" + self.record_many([ + { + "endpoint": endpoint, + "key_type": key_type, + "key": key, + "raw_freq": raw_freq, + "start_time": start_time, + "end_time": end_time, + "fields_hash": fields_hash, + "row_count": row_count, + "status": status, + "fetched_at": fetched_at, + } + ]) + + def record_many(self, rows: list[dict]) -> None: + """Append a batch of coverage rows with one normalize + one parquet write. + + Empty input is a no-op. Each row is normalized exactly like a single + ``record``; the written parquet is reindexed to ``INTRADAY_LEDGER_COLUMNS``. + No token / secret is ever written here. The in-process cache is refreshed to + the just-written frame. + """ + rows = list(rows) + if not rows: + return + normalized = [self._normalize_row(r) for r in rows] + existing = self._load() + new_rows = pd.DataFrame(normalized) + combined = new_rows if existing.empty else pd.concat( + [existing, new_rows], ignore_index=True ) combined = combined.reindex(columns=INTRADAY_LEDGER_COLUMNS) self.path.parent.mkdir(parents=True, exist_ok=True) tmp = self.path.with_suffix(".parquet.tmp") combined.to_parquet(tmp, engine="pyarrow", index=False) os.replace(tmp, self.path) + self._frame = combined + self._frame_mtime = self._current_mtime() + self._day_memo.clear() diff --git a/data/cache/tushare_cache.py b/data/cache/tushare_cache.py index fd876df..74aae54 100644 --- a/data/cache/tushare_cache.py +++ b/data/cache/tushare_cache.py @@ -482,41 +482,51 @@ def _record_gap_coverage( """ pending_start = self._pending_start() if pending_start is None or gap_end < pending_start: - self._record_one(endpoint, symbol, gap_start, gap_end, parsed, - "empty", fields_hash) + self._record_rows( + endpoint, symbol, + [(gap_start, gap_end, parsed, "empty")], fields_hash, + ) return has_date = "date" in parsed.columns hist_end = pending_start - pd.Timedelta(days=1) + entries = [] if gap_start <= hist_end: hist = parsed[parsed["date"] <= hist_end] if has_date else parsed.iloc[0:0] - self._record_one(endpoint, symbol, gap_start, hist_end, hist, - "empty", fields_hash) + entries.append((gap_start, hist_end, hist, "empty")) pend = parsed[parsed["date"] >= pending_start] if has_date else parsed - self._record_one(endpoint, symbol, pending_start, gap_end, pend, - "not_ready", fields_hash) + entries.append((pending_start, gap_end, pend, "not_ready")) + self._record_rows(endpoint, symbol, entries, fields_hash) - def _record_one( - self, endpoint, symbol, start, end, rows, empty_status, fields_hash - ): - """Append one coverage row; ``empty_status`` is used when ``rows`` is empty.""" - n = len(rows) - status = "ok" if n else empty_status - if status == "not_ready": - self.not_ready_counts[endpoint] = ( - self.not_ready_counts.get(endpoint, 0) + 1 - ) - self._ledger.record( - endpoint=endpoint, - key_type="symbol", - key=symbol, - start_date=start, - end_date=end, - fields_hash=fields_hash, - row_count=n, - status=status, - fetched_at=self._clock(), - source_version=self._source_version, - ) + def _record_rows(self, endpoint, symbol, entries, fields_hash): + """Append the coverage row(s) for ONE fetched gap in a single batch write. + + ``entries`` is ``[(start, end, rows, empty_status), ...]`` (one entry for a + plain gap, two when the not-ready pending tail is carved out). The status + per entry is ``ok`` when it has rows else its ``empty_status``; the not-ready + counter is bumped exactly as before. One ``record_many`` replaces the + per-row writes for the gap. + """ + batch = [] + for start, end, rows, empty_status in entries: + n = len(rows) + status = "ok" if n else empty_status + if status == "not_ready": + self.not_ready_counts[endpoint] = ( + self.not_ready_counts.get(endpoint, 0) + 1 + ) + batch.append({ + "endpoint": endpoint, + "key_type": "symbol", + "key": symbol, + "start_date": start, + "end_date": end, + "fields_hash": fields_hash, + "row_count": n, + "status": status, + "fetched_at": self._clock(), + "source_version": self._source_version, + }) + self._ledger.record_many(batch) # -- index_weight gap fetch (paged in <=90-day windows) ---------------- # def _fetch_index_gap(self, index_code, gap_start, gap_end, fetch, fields_hash): diff --git a/tests/test_coverage_ledger_scaling.py b/tests/test_coverage_ledger_scaling.py new file mode 100644 index 0000000..61b9b17 --- /dev/null +++ b/tests/test_coverage_ledger_scaling.py @@ -0,0 +1,277 @@ +"""D4 coverage-ledger scaling: record_many batch API + in-process lookup cache. + +Network-free, synthetic ledger rows. Asserts the batch path matches repeated +single records (schema/order/values), coverage semantics are unchanged +(ok/empty count; failed/not_ready do not), the in-process cache does not re-read +the full parquet for repeated identical lookups, an out-of-instance write is not +served stale, and the not-ready split writes its 1-2 rows in one batch. +""" + +from __future__ import annotations + +import pandas as pd + +from data.cache.coverage import LEDGER_COLUMNS, CoverageLedger +from data.cache.intraday_coverage import INTRADAY_LEDGER_COLUMNS, IntradayCoverageLedger +from data.cache.parquet_store import CacheParquetStore +from data.cache.tushare_cache import TushareCache + +_TS = pd.Timestamp("2026-06-13 21:00:00") + + +# --------------------------------------------------------------------------- # +# daily CoverageLedger +# --------------------------------------------------------------------------- # +def _daily_row(key, start, end, status, fetched_at=_TS, **over): + row = { + "endpoint": "market_daily", "key_type": "symbol", "key": key, + "start_date": start, "end_date": end, "fields_hash": "h", + "row_count": 1 if status in ("ok",) else 0, "status": status, + "fetched_at": fetched_at, "source_version": None, + } + row.update(over) + return row + + +def test_record_many_equals_repeated_record(tmp_path): + rows = [ + _daily_row("000001.SZ", "2024-01-02", "2024-01-03", "ok"), + _daily_row("000002.SZ", "2024-01-02", "2024-01-03", "empty"), + _daily_row("000003.SZ", None, None, "ok", fetched_at=pd.Timestamp("2026-06-14")), + ] + batch = CoverageLedger(str(tmp_path / "b")) + batch.record_many(rows) + + single = CoverageLedger(str(tmp_path / "s")) + for r in rows: + single.record(**{k: v for k, v in r.items()}) + + a = batch.read().reset_index(drop=True) + b = single.read().reset_index(drop=True) + assert list(a.columns) == LEDGER_COLUMNS + assert list(b.columns) == LEDGER_COLUMNS + pd.testing.assert_frame_equal(a, b) + + +def test_record_many_empty_is_noop(tmp_path): + led = CoverageLedger(str(tmp_path / "c")) + led.record_many([]) + assert not led.path.exists() # no file created for an empty batch + assert led.read().empty + + +def test_covered_intervals_counts_only_ok_empty(tmp_path): + led = CoverageLedger(str(tmp_path / "c")) + led.record_many([ + _daily_row("000001.SZ", "2024-01-02", "2024-01-03", "ok"), + _daily_row("000001.SZ", "2024-01-04", "2024-01-05", "empty"), + _daily_row("000001.SZ", "2024-01-06", "2024-01-07", "failed"), + _daily_row("000001.SZ", "2024-01-08", "2024-01-09", "not_ready"), + ]) + intervals = led.covered_intervals("market_daily", "000001.SZ") + got = {(s.strftime("%Y-%m-%d"), e.strftime("%Y-%m-%d")) for s, e in intervals} + assert got == {("2024-01-02", "2024-01-03"), ("2024-01-04", "2024-01-05")} + + +def test_snapshot_fetched_at_returns_latest_successful(tmp_path): + led = CoverageLedger(str(tmp_path / "c")) + led.record_many([ + _daily_row("g", None, None, "ok", fetched_at=pd.Timestamp("2026-06-10")), + _daily_row("g", None, None, "empty", fetched_at=pd.Timestamp("2026-06-12")), + _daily_row("g", None, None, "failed", fetched_at=pd.Timestamp("2026-06-20")), + _daily_row("g", None, None, "not_ready", fetched_at=pd.Timestamp("2026-06-21")), + ], ) + # latest among ok/empty only (failed/not_ready ignored) + assert led.snapshot_fetched_at("market_daily", "g") == pd.Timestamp("2026-06-12") + assert led.snapshot_fetched_at("market_daily", "absent") is None + + +# --------------------------------------------------------------------------- # +# in-process cache: repeated lookups do not re-read the parquet +# --------------------------------------------------------------------------- # +def test_repeated_lookups_do_not_reread_parquet(tmp_path): + writer = CoverageLedger(str(tmp_path / "c")) + writer.record_many([ + _daily_row("000001.SZ", "2024-01-02", "2024-01-03", "ok"), + _daily_row("g", None, None, "ok"), + ]) + # a FRESH (cold-cache) instance over the same file + reader = CoverageLedger(str(tmp_path / "c")) + calls = {"n": 0} + orig = reader._read_frame + + def counting(): + calls["n"] += 1 + return orig() + + reader._read_frame = counting + for _ in range(5): + reader.covered_intervals("market_daily", "000001.SZ") + reader.snapshot_fetched_at("market_daily", "g") + reader.read() + reader.read() + assert calls["n"] == 1 # one cold read; cache + memo serve the rest + + +def test_external_write_is_not_served_stale(tmp_path): + a = CoverageLedger(str(tmp_path / "c")) + # cold lookup over an absent file -> empty, caches "absent" + assert a.covered_intervals("market_daily", "000001.SZ") == [] + # a different instance writes a covering row + b = CoverageLedger(str(tmp_path / "c")) + b.record_many([_daily_row("000001.SZ", "2024-01-02", "2024-01-03", "ok")]) + # the first instance must see the new coverage (mtime invalidation), not the + # stale "absent" it cached on the cold lookup. + intervals = a.covered_intervals("market_daily", "000001.SZ") + assert len(intervals) == 1 + + +def test_read_returns_copy_not_internal_frame(tmp_path): + led = CoverageLedger(str(tmp_path / "c")) + led.record_many([_daily_row("000001.SZ", "2024-01-02", "2024-01-03", "ok")]) + frame = led.read() + frame.loc[0, "status"] = "MUTATED" + # mutating the returned copy must not corrupt the ledger's cached lookups + intervals = led.covered_intervals("market_daily", "000001.SZ") + assert len(intervals) == 1 + + +# --------------------------------------------------------------------------- # +# intraday IntradayCoverageLedger mirrors the pattern +# --------------------------------------------------------------------------- # +def _intra_row(key, start, end, status, fetched_at=_TS): + return { + "endpoint": "stk_mins_1min", "key_type": "symbol", "key": key, + "raw_freq": "1min", "start_time": start, "end_time": end, + "fields_hash": "h", "row_count": 1 if status == "ok" else 0, + "status": status, "fetched_at": fetched_at, + } + + +def test_intraday_record_many_equals_repeated_record(tmp_path): + rows = [ + _intra_row("000001.SZ", "2024-01-02 09:31:00", "2024-01-02 15:00:00", "ok"), + _intra_row("000002.SZ", "2024-01-03 09:31:00", "2024-01-03 15:00:00", "empty"), + ] + batch = IntradayCoverageLedger(str(tmp_path / "b")) + batch.record_many(rows) + single = IntradayCoverageLedger(str(tmp_path / "s")) + for r in rows: + single.record(**r) + a = batch.read().reset_index(drop=True) + b = single.read().reset_index(drop=True) + assert list(a.columns) == INTRADAY_LEDGER_COLUMNS + pd.testing.assert_frame_equal(a, b) + + +def test_intraday_covered_day_intervals_semantics_and_cache(tmp_path): + writer = IntradayCoverageLedger(str(tmp_path / "c")) + writer.record_many([ + _intra_row("000001.SZ", "2024-01-02 09:31:00", "2024-01-02 15:00:00", "ok"), + _intra_row("000001.SZ", "2024-01-03 09:31:00", "2024-01-03 15:00:00", "failed"), + ]) + reader = IntradayCoverageLedger(str(tmp_path / "c")) + calls = {"n": 0} + orig = reader._read_frame + + def counting(): + calls["n"] += 1 + return orig() + + reader._read_frame = counting + days = None + for _ in range(4): + days = reader.covered_day_intervals("stk_mins_1min", "000001.SZ", "1min") + assert calls["n"] == 1 # one cold read served all four lookups + # only the ok day counts (failed ignored), day-normalized + assert days == [(pd.Timestamp("2024-01-02"), pd.Timestamp("2024-01-02"))] + + +def test_intraday_record_many_empty_is_noop(tmp_path): + led = IntradayCoverageLedger(str(tmp_path / "c")) + led.record_many([]) + assert not led.path.exists() + assert led.read().empty + + +# --------------------------------------------------------------------------- # +# cache wiring: not-ready split writes its rows in ONE batch +# --------------------------------------------------------------------------- # +class _DBFetch: + def __init__(self, catalog): + self.catalog = catalog + self.calls: list[tuple] = [] + + def __call__(self, symbol, s, e): + self.calls.append((symbol, s, e)) + S, E = pd.Timestamp(s), pd.Timestamp(e) + rows = [r for r in self.catalog.get(symbol, []) if S <= pd.Timestamp(r[0]) <= E] + if not rows: + return pd.DataFrame() + return pd.DataFrame({ + "ts_code": [symbol] * len(rows), + "trade_date": [r[0].replace("-", "") for r in rows], + "pe": [r[1] for r in rows], "pb": [r[2] for r in rows], + "total_mv": [r[3] for r in rows], + }) + + +_DB_CAT = {"000001.SZ": [ + ("2024-01-02", 10.0, 1.0, 1000.0), + ("2024-01-03", 11.0, 1.1, 1100.0), + ("2024-01-04", 12.0, 1.2, 1200.0), +]} + + +def _cache(tmp_path, **kw): + root = str(tmp_path / "cache") + return TushareCache( + CacheParquetStore(root), CoverageLedger(root), + clock=lambda: pd.Timestamp("2026-06-13 21:00:00"), **kw, + ) + + +def test_not_ready_split_records_in_one_batch(tmp_path): + today = pd.Timestamp("2024-01-05") + cache = _cache(tmp_path, refresh_recent_days=0, not_ready_days=1, today=today) + batches: list[int] = [] + orig = cache._ledger.record_many + + def spy(rows): + rows = list(rows) + batches.append(len(rows)) + return orig(rows) + + cache._ledger.record_many = spy + f = _DBFetch(_DB_CAT) + cache.daily_basic(["000001.SZ"], "2024-01-02", "2024-01-05", f) + + # the single fetched gap carves hist (ok, through 01-04) + pending 01-05 + # (not_ready) and writes BOTH in ONE record_many call. + assert batches == [2] + led = cache._ledger.read() + assert (led["status"] == "not_ready").any() + assert (led["status"] == "ok").any() + + +def test_plain_gap_records_single_row_batch(tmp_path): + cache = _cache(tmp_path, refresh_recent_days=0) # not_ready_days=0 default + batches: list[int] = [] + orig = cache._ledger.record_many + + def spy(rows): + rows = list(rows) + batches.append(len(rows)) + return orig(rows) + + cache._ledger.record_many = spy + f = _DBFetch(_DB_CAT) + cache.daily_basic(["000001.SZ"], "2024-01-02", "2024-01-04", f) + assert batches == [1] # plain gap -> one coverage row, one batch + + +def test_ledger_columns_have_no_secret_fields(): + # structural: the ledger schema is endpoint metadata only (no token field). + for col in LEDGER_COLUMNS + INTRADAY_LEDGER_COLUMNS: + assert "token" not in col.lower() + assert "secret" not in col.lower()