"""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,
)