Source code for numeraire_dataset._hkm

"""Fixed-release loader for the He--Kelly--Manela paper archive.

The archive is either downloaded from the one official URL into memory or supplied through an
explicit local path.  Its byte identity, complete ZIP envelope, all five CSV schemas, and paper
sample calendars are validated before the two original-sample test-asset files can be returned.
The bundled Julia code is never executed or interpreted, and no archive member is extracted,
cached, copied, or persisted.

The author's data page permits free non-commercial use and supplies the files as-is, but does not
state a general redistribution licence.  Package distributions therefore contain loader code only.
"""

from __future__ import annotations

import csv
import hashlib
import io
import math
import os
import re
import stat
import struct
import zipfile
from dataclasses import dataclass
from datetime import UTC, datetime
from http.client import HTTPException
from pathlib import Path, PurePosixPath
from typing import Any, Literal, TypeAlias, TypedDict, cast
from urllib.parse import urlsplit
from urllib.request import HTTPRedirectHandler, Request, build_opener

import numpy as np
import pandas as pd

HKM_DATA_PAGE = "https://asaf.manela.org/data/"
HKM_ARCHIVE_URL = (
    "https://asaf.manela.org/papers/hkm/intermediarycapitalrisk/He_Kelly_Manela_Factors.zip"
)
HKM_ARCHIVE_SHA256 = "0e9bac2d35c80e62b20fc71f222a41f71eb8e607963f7022da55b1bb143c3bbf"

_QUARTERLY_PAPER_MEMBER = "He_Kelly_Manela_Factors_And_Test_Assets.csv"
_MONTHLY_PAPER_MEMBER = "He_Kelly_Manela_Factors_And_Test_Assets_monthly.csv"
_DAILY_FACTOR_MEMBER = "He_Kelly_Manela_Factors_daily.csv"
_MONTHLY_FACTOR_MEMBER = "He_Kelly_Manela_Factors_monthly.csv"
_QUARTERLY_FACTOR_MEMBER = "He_Kelly_Manela_Factors_quarterly.csv"
_README_MEMBER = "readme.txt"
_JULIA_MEMBER = "He_Kelly_Manela_XS_Tests.jl"

HKM_ARCHIVE_MEMBERS = (
    _README_MEMBER,
    _QUARTERLY_PAPER_MEMBER,
    _MONTHLY_PAPER_MEMBER,
    _DAILY_FACTOR_MEMBER,
    _MONTHLY_FACTOR_MEMBER,
    _QUARTERLY_FACTOR_MEMBER,
    _JULIA_MEMBER,
)

HKM_ASSET_CLASSES = (
    "ff25",
    "us_bonds",
    "sovereign_bonds",
    "options",
    "cds",
    "commodities",
    "fx",
)

HKMFrequency: TypeAlias = Literal["quarterly", "monthly"]
HKMAssetClass: TypeAlias = Literal[
    "ff25",
    "us_bonds",
    "sovereign_bonds",
    "options",
    "cds",
    "commodities",
    "fx",
]

_ASSET_CLASS_SPECS: tuple[tuple[HKMAssetClass, str, int], ...] = (
    ("ff25", "FF25", 25),
    ("us_bonds", "US_bonds", 20),
    ("sovereign_bonds", "Sov_bonds", 6),
    ("options", "Options", 18),
    ("cds", "CDS", 20),
    ("commodities", "Commod", 23),
    ("fx", "FX", 12),
)
_ASSET_COLUMNS = tuple(
    f"{prefix}_{position:02d}"
    for _, prefix, count in _ASSET_CLASS_SPECS
    for position in range(1, count + 1)
)
_ALL_COLUMNS = tuple(f"All_{position:02d}" for position in range(1, 125))
_ALLOWED_MEMBERS = frozenset(HKM_ARCHIVE_MEMBERS)

_QUARTERLY_FACTOR_PREFIX = (
    "yyyyq",
    "intermediary_capital_ratio",
    "intermediary_leverage_ratio_squared",
    "aem_leverage_ratio",
    "intermediary_capital_risk_factor",
    "intermediary_value_weighted_investment_return",
    "aem_leverage_factor",
    "mkt_rf",
    "smb",
    "hml",
    "rf",
)
_MONTHLY_FACTOR_PREFIX = (
    "yyyymm",
    "intermediary_capital_ratio",
    "intermediary_leverage_ratio_squared",
    "intermediary_capital_risk_factor",
    "intermediary_value_weighted_investment_return",
    "mkt_rf",
    "smb",
    "hml",
    "rf",
)
_UPDATED_FACTOR_HEADER = (
    "date",
    "intermediary_capital_ratio",
    "intermediary_capital_risk_factor",
    "intermediary_value_weighted_investment_return",
    "intermediary_leverage_ratio_squared",
)
_UPDATED_SOURCE_HEADERS = {
    _DAILY_FACTOR_MEMBER: ("yyyymmdd", *_UPDATED_FACTOR_HEADER[1:]),
    _MONTHLY_FACTOR_MEMBER: ("yyyymm", *_UPDATED_FACTOR_HEADER[1:]),
    _QUARTERLY_FACTOR_MEMBER: ("yyyyq", *_UPDATED_FACTOR_HEADER[1:]),
}
_PAPER_HEADERS = {
    "quarterly": (*_QUARTERLY_FACTOR_PREFIX, *_ASSET_COLUMNS, *_ALL_COLUMNS),
    "monthly": (*_MONTHLY_FACTOR_PREFIX, *_ASSET_COLUMNS, *_ALL_COLUMNS),
}
_PAPER_MEMBERS = {
    "quarterly": _QUARTERLY_PAPER_MEMBER,
    "monthly": _MONTHLY_PAPER_MEMBER,
}
_PAPER_PERIODS = {
    "quarterly": pd.period_range("1970Q1", "2012Q4", freq="Q-DEC"),
    "monthly": pd.period_range("1970-01", "2012-12", freq="M"),
}
_UPDATED_PERIODS = {
    _MONTHLY_FACTOR_MEMBER: pd.period_range("1970-01", "2018-11", freq="M"),
    _QUARTERLY_FACTOR_MEMBER: pd.period_range("1970Q1", "2018Q3", freq="Q-DEC"),
}
_UPDATED_DAILY_ROWS = 4_766
_UPDATED_DAILY_START = pd.Timestamp("2000-01-03")
_UPDATED_DAILY_END = pd.Timestamp("2018-12-11")

_MAX_ARCHIVE_BYTES = 2_000_000
_MAX_ARCHIVE_MEMBERS = len(HKM_ARCHIVE_MEMBERS)
_MAX_MEMBER_BYTES = 1_000_000
_MAX_UNCOMPRESSED_BYTES = 2_000_000
_MAX_COMPRESSION_RATIO = 100
_MAX_CSV_RECORDS = 5_000
_MAX_CSV_COLUMNS = 300
_READ_CHUNK_BYTES = 1_048_576
_OPEN_SUPPORTS_DIR_FD = os.open in os.supports_dir_fd
_OFFICIAL_HOST = "asaf.manela.org"
_RECIPE = "hkm-paper-assets-fixed-archive-total-to-excess-decimal-v1"
_PROVENANCE_ATTR = "numeraire_dataset:provenance"
_QUARTER_SOURCE = re.compile(r"([0-9]{4})([1-4])\.0000")
_MONTH_SOURCE = re.compile(r"([0-9]{4})(0[1-9]|1[0-2])\.0000")
_QUARTER_UPDATED_SOURCE = re.compile(r"([0-9]{4})([1-4])")
_MONTH_UPDATED_SOURCE = re.compile(r"([0-9]{4})(0[1-9]|1[0-2])")
_DAY_SOURCE = re.compile(r"[0-9]{8}")
_RETURN_FIELDS = frozenset(
    {
        "intermediary_capital_risk_factor",
        "intermediary_value_weighted_investment_return",
        "aem_leverage_factor",
        "mkt_rf",
        "smb",
        "hml",
        "rf",
    }
)


@dataclass(frozen=True)
class _Release:
    url: str
    content_sha256: str
    byte_count: int


_RELEASE = _Release(
    url=HKM_ARCHIVE_URL,
    content_sha256=HKM_ARCHIVE_SHA256,
    byte_count=1_454_750,
)


class _SourceMetadata(TypedDict):
    source_mode: str
    resolved_url: str
    retrieved_at: str
    http_status: str
    http_content_type: str
    http_content_length: str


@dataclass(frozen=True)
class _ArchiveContents:
    selected: dict[str, bytes]
    member_sha256: dict[str, str]
    member_bytes: dict[str, int]
    member_manifest_sha256: str


def _restamp(frame: pd.DataFrame, **updates: str) -> pd.DataFrame:
    result = frame.copy()
    existing = result.attrs.get(_PROVENANCE_ATTR)
    if not isinstance(existing, dict) or not all(
        isinstance(key, str) and isinstance(value, str) for key, value in existing.items()
    ):
        raise ValueError("HKM frame has no valid numeraire-dataset provenance")
    provenance = cast(dict[str, str], existing).copy()
    provenance.update(updates)
    result.attrs[_PROVENANCE_ATTR] = provenance
    return result


[docs] @dataclass(frozen=True) class HKMAssetClassData: """One asset class aligned on its joint complete-case paper sample.""" frequency: HKMFrequency asset_class: HKMAssetClass factors: pd.DataFrame excess_returns: pd.DataFrame asset_metadata: pd.DataFrame
[docs] @dataclass(frozen=True) class HKMPaperData: """Fixed paper-sample factors and the unbalanced 124-test-asset panel.""" frequency: HKMFrequency factors: pd.DataFrame excess_returns: pd.DataFrame asset_metadata: pd.DataFrame
[docs] def complete_case(self, asset_class: HKMAssetClass) -> HKMAssetClassData: """Return one asset class on dates where all its assets and factors are observed.""" if asset_class not in HKM_ASSET_CLASSES: raise ValueError(f"unknown HKM asset class {asset_class!r}") metadata = self.asset_metadata.loc[ self.asset_metadata["asset_class"].eq(asset_class) ].reset_index(drop=True) assets = cast(list[str], metadata["asset"].tolist()) if not assets: raise ValueError(f"HKM asset class {asset_class!r} has no assets") if not self.factors["date"].equals(self.excess_returns["date"]): raise ValueError("HKM factor and test-asset calendars are not aligned") factor_columns = [column for column in self.factors.columns if column != "date"] complete = self.factors[factor_columns].notna().all(axis=1) complete &= self.excess_returns[assets].notna().all(axis=1) if not bool(complete.any()): raise ValueError(f"HKM asset class {asset_class!r} has no joint complete-case rows") factors = self.factors.loc[complete].reset_index(drop=True) returns = self.excess_returns.loc[complete, ["date", *assets]].reset_index(drop=True) selected_rows = str(len(factors)) selected_assets = str(len(assets)) factors = _restamp( factors, frame_role="complete_case_factors", selected_asset_class=asset_class, selected_rows=selected_rows, selected_assets=selected_assets, ) returns = _restamp( returns, frame_role="complete_case_excess_returns", selected_asset_class=asset_class, selected_rows=selected_rows, selected_assets=selected_assets, ) metadata = _restamp( metadata, frame_role="complete_case_asset_metadata", selected_asset_class=asset_class, selected_rows=selected_rows, selected_assets=selected_assets, ) return HKMAssetClassData( frequency=self.frequency, asset_class=asset_class, factors=factors, excess_returns=returns, asset_metadata=metadata, )
def _validate_official_url(url: str) -> None: parsed = urlsplit(url) try: port = parsed.port except ValueError as exc: raise ValueError("HKM download resolved to an untrusted URL") from exc if ( url != HKM_ARCHIVE_URL or parsed.scheme != "https" or parsed.hostname != _OFFICIAL_HOST or parsed.username is not None or parsed.password is not None or port is not None or parsed.query or parsed.fragment ): raise ValueError(f"HKM download resolved to an untrusted URL {url!r}") def _header(headers: object, name: str) -> str | None: getter = getattr(headers, "get", None) if not callable(getter): return None value = getter(name) if value is None: return None if not isinstance(value, str): raise ValueError(f"HKM download returned a non-text {name} header") return value.strip() class _OfficialRedirectHandler(HTTPRedirectHandler): """Reject every redirect away from the one frozen official archive URL.""" def redirect_request( self, req: Request, fp: Any, code: int, msg: str, headers: Any, newurl: str, ) -> Request | None: _validate_official_url(newurl) return super().redirect_request(req, fp, code, msg, headers, newurl) def _open_official(request: Request, *, timeout: float) -> Any: return build_opener(_OfficialRedirectHandler()).open(request, timeout=timeout) def _download_official(*, timeout: float) -> tuple[bytes, _SourceMetadata]: _validate_official_url(_RELEASE.url) if not math.isfinite(timeout) or timeout <= 0: raise ValueError("timeout must be positive and finite") request = Request( _RELEASE.url, headers={"User-Agent": "numeraire-dataset (public academic data loader)"}, ) try: with _open_official(request, timeout=timeout) as response: resolved_url = response.geturl() _validate_official_url(resolved_url) getcode = getattr(response, "getcode", None) status = getcode() if callable(getcode) else None if status is not None and status != 200: raise ValueError(f"HKM download returned HTTP status {status}") content = response.read(_MAX_ARCHIVE_BYTES + 1) content_type = _header(response.headers, "Content-Type") content_length = _header(response.headers, "Content-Length") retrieved_at = datetime.now(UTC).isoformat().replace("+00:00", "Z") except (OSError, HTTPException) as exc: raise RuntimeError("could not download the official HKM paper archive") from exc if not content: raise ValueError("HKM archive download was empty") if len(content) > _MAX_ARCHIVE_BYTES: raise ValueError("HKM archive download exceeded the compressed-size safety limit") if content_length is not None: try: declared_length = int(content_length) except ValueError as exc: raise ValueError("HKM download returned an invalid Content-Length header") from exc if declared_length != len(content): raise ValueError("HKM Content-Length does not match response bytes") return content, { "source_mode": "official_memory_download", "resolved_url": resolved_url, "retrieved_at": retrieved_at, "http_status": str(status) if status is not None else "unreported", "http_content_type": content_type or "unreported", "http_content_length": content_length or "unreported", } def _local_bytes(path: str | Path) -> tuple[bytes, _SourceMetadata]: try: source = Path(path).expanduser() except TypeError as exc: raise ValueError("HKM archive path must be a string or pathlib.Path") from exc if ".." in source.parts: raise ValueError("HKM archive path must not contain '..' components") absolute = Path(os.path.abspath(source)) if absolute.suffix.casefold() != ".zip": raise ValueError("HKM input must be a .zip archive") required_flags = ("O_DIRECTORY", "O_NOFOLLOW", "O_NONBLOCK") if ( os.name != "posix" or not _OPEN_SUPPORTS_DIR_FD or any(not hasattr(os, flag) for flag in required_flags) or absolute.anchor != os.sep ): raise ValueError("secure HKM local-file opening is unsupported on this platform") close_on_exec = getattr(os, "O_CLOEXEC", 0) directory_flags = os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW | close_on_exec file_flags = os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK | close_on_exec try: directory_fd = os.open(absolute.anchor, directory_flags) except OSError as exc: raise ValueError( "HKM archive path components must be readable non-symlink directories" ) from exc try: for component in absolute.parts[1:-1]: try: next_fd = os.open(component, directory_flags, dir_fd=directory_fd) except OSError as exc: raise ValueError( "HKM archive path components must be readable non-symlink directories" ) from exc os.close(directory_fd) directory_fd = next_fd try: file_fd = os.open(absolute.name, file_flags, dir_fd=directory_fd) except OSError as exc: raise ValueError("HKM archive path must be a readable non-symlink local file") from exc finally: os.close(directory_fd) try: before = os.fstat(file_fd) if not stat.S_ISREG(before.st_mode): raise ValueError("HKM archive path is not a regular file") if before.st_size <= 0: raise ValueError("HKM archive is empty") if before.st_size > _MAX_ARCHIVE_BYTES: raise ValueError("HKM archive exceeds the compressed-size safety limit") content = bytearray() while len(content) <= _MAX_ARCHIVE_BYTES: remaining = _MAX_ARCHIVE_BYTES + 1 - len(content) chunk = os.read(file_fd, min(_READ_CHUNK_BYTES, remaining)) if not chunk: break content.extend(chunk) after = os.fstat(file_fd) except OSError as exc: raise ValueError("HKM archive path is not readable") from exc finally: os.close(file_fd) if len(content) > _MAX_ARCHIVE_BYTES: raise ValueError("HKM archive exceeds the compressed-size safety limit") before_signature = ( before.st_dev, before.st_ino, before.st_mode, before.st_size, before.st_mtime_ns, before.st_ctime_ns, ) after_signature = ( after.st_dev, after.st_ino, after.st_mode, after.st_size, after.st_mtime_ns, after.st_ctime_ns, ) if len(content) != before.st_size or before_signature != after_signature: raise ValueError("HKM archive changed while it was being read") return bytes(content), { "source_mode": "caller_path_exact_official_archive", "resolved_url": HKM_ARCHIVE_URL, "retrieved_at": "caller_supplied", "http_status": "not_applicable", "http_content_type": "not_applicable", "http_content_length": "not_applicable", } def _validate_release(content: bytes) -> str: if len(content) != _RELEASE.byte_count: raise ValueError("HKM archive does not match the frozen official byte count") digest = hashlib.sha256(content).hexdigest() if digest != _RELEASE.content_sha256: raise ValueError("HKM archive SHA-256 does not match the frozen official release") return digest def _safe_member_name(name: str) -> str: if not name or "\\" in name or "\0" in name: raise ValueError("HKM archive contains an unsafe member name") path = PurePosixPath(name) if path.is_absolute() or any(part in {"", ".", ".."} for part in path.parts): raise ValueError("HKM archive contains an unsafe member path") if len(path.parts) != 1 or path.as_posix() != name: raise ValueError("HKM archive members must be canonical root-level files") return name def _validate_zip_container(content: bytes) -> None: """Require one conventional ZIP ending exactly at its comment-free EOCD record.""" signature = b"PK\x05\x06" minimum_size = struct.calcsize("<4s4H2LH") if len(content) < minimum_size: raise ValueError("HKM input is not a complete ZIP archive") offset = content.rfind(signature, max(0, len(content) - 65_557)) if offset < 0 or offset + minimum_size != len(content): raise ValueError("HKM ZIP contains a trailing or prepended payload") ( observed_signature, disk_number, central_disk, disk_entries, total_entries, central_size, central_offset, comment_length, ) = struct.unpack("<4s4H2LH", content[offset:]) if ( observed_signature != signature or disk_number != 0 or central_disk != 0 or disk_entries != total_entries or total_entries == 0 or total_entries > _MAX_ARCHIVE_MEMBERS or comment_length != 0 or central_offset + central_size != offset ): raise ValueError("HKM ZIP has an unsupported central-directory envelope") def _read_archive(content: bytes) -> _ArchiveContents: _validate_zip_container(content) selected: dict[str, bytes] = {} digests: dict[str, str] = {} sizes: dict[str, int] = {} try: with zipfile.ZipFile(io.BytesIO(content), mode="r") as archive: if archive.comment: raise ValueError("HKM archive must not contain a ZIP comment") members = archive.infolist() if len(members) != _MAX_ARCHIVE_MEMBERS: raise ValueError("HKM archive must contain exactly the documented member set") seen: set[str] = set() total = 0 validated: list[tuple[zipfile.ZipInfo, str]] = [] for member in members: name = _safe_member_name(member.filename) folded = name.casefold() if folded in seen: raise ValueError("HKM archive has duplicate member names") seen.add(folded) if name not in _ALLOWED_MEMBERS: raise ValueError(f"HKM archive contains unexpected member {name!r}") if member.is_dir(): raise ValueError("HKM archive must contain only regular files") mode = (member.external_attr >> 16) & 0xFFFF if member.create_system == 3 and mode and not stat.S_ISREG(mode): raise ValueError("HKM archive must contain only regular files, never links") if member.flag_bits & 0x1: raise ValueError("HKM archive members must not be encrypted") if member.comment or member.extra: raise ValueError("HKM archive members must not contain hidden metadata fields") if member.compress_type not in {zipfile.ZIP_STORED, zipfile.ZIP_DEFLATED}: raise ValueError("HKM archive uses an unsupported compression method") if member.file_size <= 0 or member.file_size > _MAX_MEMBER_BYTES: raise ValueError("HKM archive member has an invalid or unsafe size") if member.compress_size <= 0: raise ValueError("HKM archive member has an invalid compression envelope") if member.file_size / member.compress_size > _MAX_COMPRESSION_RATIO: raise ValueError("HKM archive member has an unsafe compression ratio") total += member.file_size if total > _MAX_UNCOMPRESSED_BYTES: raise ValueError("HKM archive exceeds the uncompressed-size safety limit") validated.append((member, name)) if seen != {name.casefold() for name in _ALLOWED_MEMBERS}: raise ValueError("HKM archive must contain exactly the documented member set") if total / len(content) > _MAX_COMPRESSION_RATIO: raise ValueError("HKM archive has an unsafe aggregate compression ratio") for member, name in validated: digest = hashlib.sha256() chunks: list[bytes] = [] count = 0 with archive.open(member, mode="r") as stream: while True: chunk = stream.read(_READ_CHUNK_BYTES) if not chunk: break count += len(chunk) if count > member.file_size or count > _MAX_MEMBER_BYTES: raise ValueError( f"HKM archive member {name!r} exceeds its declared size" ) digest.update(chunk) chunks.append(chunk) if count != member.file_size: raise ValueError(f"HKM archive member {name!r} is truncated") selected[name] = b"".join(chunks) digests[name] = digest.hexdigest() sizes[name] = count except (zipfile.BadZipFile, zipfile.LargeZipFile, OSError, RuntimeError) as exc: raise ValueError("HKM input is not a valid safe ZIP archive") from exc manifest = hashlib.sha256( "\n".join(f"{name}:{digests[name]}:{sizes[name]}" for name in sorted(digests)).encode( "ascii" ) ).hexdigest() return _ArchiveContents( selected=selected, member_sha256=digests, member_bytes=sizes, member_manifest_sha256=manifest, ) def _csv_rows(content: bytes, *, member: str) -> list[list[str]]: try: text = content.decode("ascii") except UnicodeDecodeError as exc: raise ValueError(f"HKM member {member!r} is not ASCII CSV") from exc if "\0" in text or '"' in text: raise ValueError(f"HKM member {member!r} contains unsupported CSV syntax") physical_lines = text.splitlines() if not physical_lines or len(physical_lines) > _MAX_CSV_RECORDS: raise ValueError(f"HKM member {member!r} has an unsafe CSV record count") if any(not line for line in physical_lines): raise ValueError(f"HKM member {member!r} contains a blank CSV record") try: rows = list(csv.reader(physical_lines, strict=True)) except csv.Error as exc: raise ValueError(f"HKM member {member!r} is malformed CSV") from exc if len(rows) != len(physical_lines): raise ValueError(f"HKM member {member!r} contains multiline CSV records") if any(len(row) > _MAX_CSV_COLUMNS for row in rows): raise ValueError(f"HKM member {member!r} exceeds the CSV column-count limit") if any(cell != cell.strip(" \t\r\n\v\f") for row in rows for cell in row): raise ValueError(f"HKM member {member!r} contains whitespace-padded fields") return rows def _number(cell: str, *, member: str, field: str, allow_missing: bool = False) -> float: if not cell: if allow_missing: return float("nan") raise ValueError(f"HKM member {member!r} has a missing {field}") try: value = float(cell) except ValueError as exc: raise ValueError(f"HKM member {member!r} has a non-numeric {field}") from exc if not math.isfinite(value): raise ValueError(f"HKM member {member!r} has a non-finite {field}") if field == "intermediary_capital_ratio" and not 0.0 <= value <= 1.0: raise ValueError(f"HKM member {member!r} violates the capital-ratio envelope") if field in {"intermediary_leverage_ratio_squared", "aem_leverage_ratio"} and not ( 0.0 < value <= 100_000.0 ): raise ValueError(f"HKM member {member!r} violates the leverage-ratio envelope") if (field in _RETURN_FIELDS or field in _ASSET_COLUMNS) and abs(value) > 2.0: raise ValueError(f"HKM member {member!r} violates the decimal-return unit envelope") return value def _paper_period(cell: str, *, frequency: HKMFrequency, member: str) -> pd.Period: pattern = _QUARTER_SOURCE if frequency == "quarterly" else _MONTH_SOURCE match = pattern.fullmatch(cell) if match is None: raise ValueError(f"HKM member {member!r} has an invalid {frequency} date") year, subperiod = (int(value) for value in match.groups()) if frequency == "quarterly": return pd.Period(year=year, quarter=subperiod, freq="Q-DEC") return pd.Period(year=year, month=subperiod, freq="M") def _updated_period(cell: str, *, frequency: str, member: str) -> pd.Period | pd.Timestamp: if frequency == "daily": if _DAY_SOURCE.fullmatch(cell) is None: raise ValueError(f"HKM member {member!r} has an invalid daily date") try: stamp = pd.Timestamp(datetime.strptime(cell, "%Y%m%d")) except ValueError as exc: raise ValueError(f"HKM member {member!r} has an invalid daily date") from exc if stamp.strftime("%Y%m%d") != cell: raise ValueError(f"HKM member {member!r} has a non-canonical daily date") return stamp pattern = _MONTH_UPDATED_SOURCE if frequency == "monthly" else _QUARTER_UPDATED_SOURCE match = pattern.fullmatch(cell) if match is None: raise ValueError(f"HKM member {member!r} has an invalid {frequency} date") year, subperiod = (int(value) for value in match.groups()) if frequency == "monthly": return pd.Period(year=year, month=subperiod, freq="M") return pd.Period(year=year, quarter=subperiod, freq="Q-DEC") def _parse_paper_member( content: bytes, *, frequency: HKMFrequency, ) -> tuple[pd.DataFrame, pd.DataFrame]: member = _PAPER_MEMBERS[frequency] rows = _csv_rows(content, member=member) expected_header = _PAPER_HEADERS[frequency] if tuple(rows[0]) != expected_header: raise ValueError(f"HKM member {member!r} has an unexpected paper schema") expected_periods = _PAPER_PERIODS[frequency] if len(rows) != len(expected_periods) + 1: raise ValueError(f"HKM member {member!r} has an unexpected paper row count") width = len(expected_header) if any(len(row) != width for row in rows[1:]): raise ValueError(f"HKM member {member!r} has inconsistent CSV row widths") prefix = _QUARTERLY_FACTOR_PREFIX if frequency == "quarterly" else _MONTHLY_FACTOR_PREFIX prefix_width = len(prefix) periods: list[pd.Period] = [] factor_values: dict[str, list[float]] = {field: [] for field in prefix[1:]} asset_values = np.empty((len(expected_periods), len(_ASSET_COLUMNS)), dtype=np.float64) for row_number, row in enumerate(rows[1:], start=1): period = _paper_period(row[0], frequency=frequency, member=member) periods.append(period) base_cells = row[prefix_width : prefix_width + len(_ASSET_COLUMNS)] all_cells = row[prefix_width + len(_ASSET_COLUMNS) :] if base_cells != all_cells: raise ValueError( f"HKM member {member!r} All_01..All_124 do not exactly duplicate test assets" ) for position, field in enumerate(prefix[1:], start=1): factor_values[field].append(_number(row[position], member=member, field=field)) for column, (field, cell) in enumerate(zip(_ASSET_COLUMNS, base_cells, strict=True)): asset_values[row_number - 1, column] = _number( cell, member=member, field=field, allow_missing=True, ) observed = pd.PeriodIndex(periods, freq=expected_periods.freq) if observed.has_duplicates or not observed.equals(expected_periods): raise ValueError(f"HKM member {member!r} has an unexpected or duplicate paper calendar") if np.isnan(asset_values).all(axis=0).any(): raise ValueError(f"HKM member {member!r} has an entirely missing test asset") dates = observed.to_timestamp(how="end").normalize() factors = pd.DataFrame({"date": dates}) canonical_factor_names = { "mkt_rf": "mkt_excess", "rf": "risk_free", } for field in prefix[1:]: factors[canonical_factor_names.get(field, field)] = factor_values[field] assets = pd.DataFrame(asset_values, columns=_ASSET_COLUMNS) risk_free = factors["risk_free"].to_numpy(dtype=np.float64) assets = assets.subtract(risk_free, axis=0) assets.insert(0, "date", dates) return factors, assets def _validate_updated_member(content: bytes, *, member: str, frequency: str) -> None: rows = _csv_rows(content, member=member) if tuple(rows[0]) != _UPDATED_SOURCE_HEADERS[member]: raise ValueError(f"HKM member {member!r} has an unexpected factor-only schema") if any(len(row) != len(_UPDATED_FACTOR_HEADER) for row in rows[1:]): raise ValueError(f"HKM member {member!r} has inconsistent CSV row widths") dates: list[pd.Period | pd.Timestamp] = [] for row_number, row in enumerate(rows[1:], start=1): dates.append(_updated_period(row[0], frequency=frequency, member=member)) for position, field in enumerate(_UPDATED_FACTOR_HEADER[1:], start=1): allow_missing = ( member == _DAILY_FACTOR_MEMBER and row_number == 1 and field == "intermediary_capital_risk_factor" ) if allow_missing and row[position]: raise ValueError(f"HKM member {member!r} violates the frozen missing-data envelope") _number(row[position], member=member, field=field, allow_missing=allow_missing) if frequency == "daily": stamps = pd.DatetimeIndex(cast(list[pd.Timestamp], dates)) if ( len(stamps) != _UPDATED_DAILY_ROWS or stamps.has_duplicates or not stamps.is_monotonic_increasing or stamps[0] != _UPDATED_DAILY_START or stamps[-1] != _UPDATED_DAILY_END or bool((stamps.dayofweek > 4).any()) ): raise ValueError(f"HKM member {member!r} has an unexpected daily calendar") return expected = _UPDATED_PERIODS[member] periods = pd.PeriodIndex(cast(list[pd.Period], dates), freq=expected.freq) if periods.has_duplicates or not periods.equals(expected): raise ValueError(f"HKM member {member!r} has an unexpected {frequency} calendar") def _asset_metadata() -> pd.DataFrame: rows: list[dict[str, str | int]] = [] global_position = 1 for asset_class, prefix, count in _ASSET_CLASS_SPECS: for class_position in range(1, count + 1): rows.append( { "asset": f"{prefix}_{class_position:02d}", "asset_class": asset_class, "class_position": class_position, "source_all_column": f"All_{global_position:02d}", } ) global_position += 1 return pd.DataFrame(rows) def _provenance( *, frequency: HKMFrequency, archive_sha256: str, contents: _ArchiveContents, source_metadata: _SourceMetadata, factors: pd.DataFrame, ) -> dict[str, str]: member = _PAPER_MEMBERS[frequency] return { "backend": "hkm-official-fixed-archive", "dataset": "He-Kelly-Manela Intermediary Asset Pricing", "source_page": HKM_DATA_PAGE, "source_url": HKM_ARCHIVE_URL, "archive_sha256": archive_sha256, "archive_bytes": str(_RELEASE.byte_count), "archive_members": str(len(contents.member_sha256)), "archive_uncompressed_bytes": str(sum(contents.member_bytes.values())), "member_manifest_sha256": contents.member_manifest_sha256, "member_name": member, "member_sha256": contents.member_sha256[member], "member_bytes": str(contents.member_bytes[member]), "frequency": frequency, "paper_sample": ("1970Q1-2012Q4" if frequency == "quarterly" else "1970-01-2012-12"), "full_start_date": pd.Timestamp(factors["date"].iloc[0]).date().isoformat(), "full_end_date": pd.Timestamp(factors["date"].iloc[-1]).date().isoformat(), "full_rows": str(len(factors)), "asset_count": str(len(_ASSET_COLUMNS)), "asset_classes": ",".join(HKM_ASSET_CLASSES), "source_return_unit": "decimal_total_or_net_return", "output_return_unit": "decimal_excess_return", "source_factor_unit": "decimal", "risk_free_subtraction": "test_asset_minus_same_period_rf", "all_columns_validation": "exact_raw_cell_duplication_verified_not_returned", "unbalanced_panel": "preserved", "updated_factor_only_members": "validated_not_exposed", "identity_validation": "expected_archive_sha256_verified", "recipe": _RECIPE, "vintage_semantics": "fixed_original_paper_archive_non_pit", "license_caveat": "free_noncommercial_use_as_is_no_general_redistribution_grant", "third_party_test_asset_sources": "cite_original_sources", "redistributable": "false", "data_vintage": f"hkm-paper-archive@sha256:{archive_sha256}/recipe:{_RECIPE}", "source_mode": source_metadata["source_mode"], "resolved_url": source_metadata["resolved_url"], "retrieved_at": source_metadata["retrieved_at"], "http_status": source_metadata["http_status"], "http_content_type": source_metadata["http_content_type"], "http_content_length": source_metadata["http_content_length"], } def load_paper_data( *, frequency: HKMFrequency, path: str | Path | None, timeout: float, ) -> tuple[HKMPaperData, dict[str, str]]: """Validate the fixed archive and return one original-sample paper frequency.""" if frequency not in {"quarterly", "monthly"}: raise ValueError("frequency must be exactly 'quarterly' or 'monthly'") if path is None: content, source_metadata = _download_official(timeout=timeout) else: content, source_metadata = _local_bytes(path) archive_sha256 = _validate_release(content) contents = _read_archive(content) _validate_updated_member( contents.selected[_DAILY_FACTOR_MEMBER], member=_DAILY_FACTOR_MEMBER, frequency="daily", ) _validate_updated_member( contents.selected[_MONTHLY_FACTOR_MEMBER], member=_MONTHLY_FACTOR_MEMBER, frequency="monthly", ) _validate_updated_member( contents.selected[_QUARTERLY_FACTOR_MEMBER], member=_QUARTERLY_FACTOR_MEMBER, frequency="quarterly", ) parsed = { selected_frequency: _parse_paper_member( contents.selected[_PAPER_MEMBERS[selected_frequency]], frequency=cast(HKMFrequency, selected_frequency), ) for selected_frequency in ("quarterly", "monthly") } factors, excess_returns = parsed[frequency] metadata = _asset_metadata() provenance = _provenance( frequency=frequency, archive_sha256=archive_sha256, contents=contents, source_metadata=source_metadata, factors=factors, ) return ( HKMPaperData( frequency=frequency, factors=factors, excess_returns=excess_returns, asset_metadata=metadata, ), provenance, )