Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
157 changes: 120 additions & 37 deletions data/cache/coverage.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,32 +52,71 @@


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)``.

Only ``ok`` / ``empty`` rows count (a ``failed`` fetch is not coverage).
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)
Expand All @@ -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)``.
Expand All @@ -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,
*,
Expand All @@ -128,30 +191,50 @@ 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)
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)
# 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()
120 changes: 96 additions & 24 deletions data/cache/intraday_coverage.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]:
Expand All @@ -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)
Expand All @@ -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,
*,
Expand All @@ -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()
Loading