Source code for numeraire_dataset._daniel_moskowitz

"""Strict path-only parser for the Daniel--Moskowitz momentum archive.

The caller supplies a locally obtained ``DM_data_2014_02.tar.gz`` path and an explicit parser
contract.  This module never downloads, discovers, extracts, copies, caches, or persists archive
bytes.  The author's public documentation identifies the archive members and their portfolio
semantics, but does not specify a machine-readable text layout.  Delimiters, header/data records,
source columns, date formats, units, and optional release envelopes therefore remain caller-pinned
contract fields rather than guessed defaults.

Secure path traversal uses POSIX ``dir_fd``/``O_NOFOLLOW`` primitives and fails closed before
reading on platforms that do not provide them.
"""

from __future__ import annotations

import codecs
import csv
import gzip
import hashlib
import io
import json
import math
import os
import re
import stat
import tarfile
from collections.abc import Iterable, Iterator
from dataclasses import asdict, dataclass
from datetime import datetime
from pathlib import Path, PurePosixPath
from typing import Literal, TypeAlias

import numpy as np
import pandas as pd

DANIEL_MOSKOWITZ_DATA_PAGE = "https://www.kentdaniel.net/data.php"
DANIEL_MOSKOWITZ_DOCUMENTATION_URL = "https://www.kentdaniel.net/data/momentum/mom_data.pdf"
DANIEL_MOSKOWITZ_ARCHIVE_FILENAME = "DM_data_2014_02.tar.gz"
DANIEL_MOSKOWITZ_DAILY_TOTAL_MEMBER = "d_m_pt_tot.txt"
DANIEL_MOSKOWITZ_MONTHLY_TOTAL_MEMBER = "m_m_pt_tot.txt"

# Exact root-level filenames in Table 1 of the author's portfolio documentation.  A tuple makes
# the public allowlist immutable; the archive parser rejects every name outside it.
DANIEL_MOSKOWITZ_ARCHIVE_MEMBERS = (
    "d_m_pt_ind.txt",
    "d_m_pt_res.txt",
    DANIEL_MOSKOWITZ_DAILY_TOTAL_MEMBER,
    "d_m_pt_nyse_ind.txt",
    "d_m_pt_nyse_res.txt",
    "d_m_pt_nyse_tot.txt",
    "m_m_pt_ind.txt",
    "m_m_pt_res.txt",
    DANIEL_MOSKOWITZ_MONTHLY_TOTAL_MEMBER,
    "m_m_pt_nyse_ind.txt",
    "m_m_pt_nyse_res.txt",
    "m_m_pt_nyse_tot.txt",
)

_ALLOWED_MEMBERS = frozenset(DANIEL_MOSKOWITZ_ARCHIVE_MEMBERS)
_MAX_ARCHIVE_BYTES = 100_000_000
_MAX_MEMBER_BYTES = 50_000_000
_MAX_UNCOMPRESSED_BYTES = 300_000_000
_MAX_COMPRESSION_RATIO = 1_000
_MAX_MEMBERS = len(DANIEL_MOSKOWITZ_ARCHIVE_MEMBERS)
_READ_CHUNK_BYTES = 1_048_576
_MAX_DECIMAL_LONG_ONLY_RETURN = 10.0
_MAX_TEXT_PHYSICAL_LINES = 150_000
_MAX_TEXT_RECORDS = 100_000
_MAX_TEXT_COLUMNS = 128
_MAX_TEXT_FIELDS = 1_000_000
_MAX_TEXT_RECORD_CHARS = 100_000
_MAX_TAR_OVERHEAD_BYTES = 2_000_000
_MAX_TAR_TRAILING_ZERO_BYTES = 1_000_000
_ASCII_WHITESPACE = " \t\n\r\v\f"
_ASCII_WHITESPACE_PATTERN = re.compile(r"[ \t\n\r\v\f]+")
_OPEN_SUPPORTS_DIR_FD = os.open in os.supports_dir_fd
_RECIPE = "dm-2014-02-explicit-text-contract-decimal-simple-10minus1-v1"
_HEX = frozenset("0123456789abcdef")

DanielMoskowitzColumn: TypeAlias = str | int
DanielMoskowitzSourceUnit: TypeAlias = Literal["decimal", "percent"]


def _column_reference(value: object, *, label: str) -> None:
    if isinstance(value, bool) or not isinstance(value, (str, int)):
        raise ValueError(f"{label} must be a string header or integer column position")
    if isinstance(value, str) and not value:
        raise ValueError(f"{label} must not be empty")
    if isinstance(value, int) and value < 0:
        raise ValueError(f"{label} integer positions must be nonnegative")


def _optional_record(value: object, *, label: str) -> None:
    if value is not None and (isinstance(value, bool) or not isinstance(value, int) or value < 0):
        raise ValueError(f"{label} must be a nonnegative integer or None")


def _optional_digest(value: object, *, label: str) -> None:
    if value is None:
        return
    if not isinstance(value, str) or len(value) != 64 or not set(value.casefold()) <= _HEX:
        raise ValueError(f"{label} must be a 64-character SHA-256 hex digest or None")


[docs] @dataclass(frozen=True) class DanielMoskowitzMemberContract: """Explicit parser and release envelope for one documented archive member. ``header_row``, ``data_start_row``, and ``data_end_row`` are zero-based positions among non-empty parsed records; ``data_end_row`` is exclusive. With ``header_row=None``, column references must be integer positions. ``delimiter="whitespace"`` means one or more ASCII whitespace characters; every other accepted delimiter is one literal character. No source-unit default exists: callers must declare either source decimals or percentage points. The loader always returns decimal simple returns. Optional SHA/date/row envelopes turn an otherwise structural parse into a caller-pinned release check. """ member_name: str delimiter: str header_row: int | None date_column: DanielMoskowitzColumn portfolio_columns: tuple[DanielMoskowitzColumn, ...] source_unit: DanielMoskowitzSourceUnit date_format: str data_start_row: int | None = None data_end_row: int | None = None expected_start_date: str | None = None expected_end_date: str | None = None expected_rows: int | None = None expected_dates: tuple[str, ...] | None = None expected_sha256: str | None = None encoding: str = "ascii" def __post_init__(self) -> None: if self.member_name not in _ALLOWED_MEMBERS: raise ValueError("member_name is not in the documented Daniel--Moskowitz allowlist") if not isinstance(self.delimiter, str) or not self.delimiter: raise ValueError("delimiter must be one literal character or 'whitespace'") if self.delimiter != "whitespace" and ( len(self.delimiter) != 1 or self.delimiter in {"\r", "\n", "\0"} ): raise ValueError("delimiter must be one literal character or 'whitespace'") _optional_record(self.header_row, label="header_row") _optional_record(self.data_start_row, label="data_start_row") _optional_record(self.data_end_row, label="data_end_row") if self.header_row is None and ( isinstance(self.date_column, str) or any(isinstance(column, str) for column in self.portfolio_columns) ): raise ValueError("string column references require a header_row") _column_reference(self.date_column, label="date_column") if not isinstance(self.portfolio_columns, tuple) or len(self.portfolio_columns) != 10: raise ValueError("portfolio_columns must be a tuple of exactly ten source columns") for index, column in enumerate(self.portfolio_columns, start=1): _column_reference(column, label=f"portfolio_columns[{index}]") if len(set(self.portfolio_columns)) != 10: raise ValueError("portfolio_columns must not contain duplicate references") if self.date_column in self.portfolio_columns: raise ValueError("date_column must be distinct from portfolio_columns") if self.source_unit not in {"decimal", "percent"}: raise ValueError("source_unit must be exactly 'decimal' or 'percent'") if not isinstance(self.date_format, str) or not self.date_format: raise ValueError("date_format must be a non-empty explicit strptime format") start = self.data_start_row if start is not None and self.header_row is not None and start <= self.header_row: raise ValueError("data_start_row must follow header_row") if self.data_end_row is not None: implied_start = 0 if start is None and self.header_row is None else start if implied_start is None: implied_start = self.header_row + 1 if self.header_row is not None else 0 if self.data_end_row <= implied_start: raise ValueError("data_end_row must be after the first data record") for value, label in ( (self.expected_start_date, "expected_start_date"), (self.expected_end_date, "expected_end_date"), ): if value is not None and (not isinstance(value, str) or not value): raise ValueError(f"{label} must be a non-empty ISO-like date or None") if self.expected_rows is not None and ( isinstance(self.expected_rows, bool) or not isinstance(self.expected_rows, int) or self.expected_rows <= 0 ): raise ValueError("expected_rows must be a positive integer or None") if self.expected_dates is not None and ( not isinstance(self.expected_dates, tuple) or not self.expected_dates or any(not isinstance(value, str) or not value for value in self.expected_dates) ): raise ValueError("expected_dates must be a non-empty tuple of date strings or None") _optional_digest(self.expected_sha256, label="expected_sha256") if not isinstance(self.encoding, str) or not self.encoding: raise ValueError("encoding must be a non-empty codec name") try: codecs.lookup(self.encoding) except LookupError as exc: raise ValueError(f"unknown text encoding {self.encoding!r}") from exc
[docs] @dataclass(frozen=True) class DanielMoskowitzArchiveContract: """Paired daily/monthly contract for the total-return, all-firm portfolios.""" daily: DanielMoskowitzMemberContract monthly: DanielMoskowitzMemberContract expected_archive_sha256: str | None = None daily_monthly_alignment_atol: float | None = None def __post_init__(self) -> None: if not isinstance(self.daily, DanielMoskowitzMemberContract) or not isinstance( self.monthly, DanielMoskowitzMemberContract ): raise ValueError("daily and monthly must be DanielMoskowitzMemberContract values") if self.daily.member_name != DANIEL_MOSKOWITZ_DAILY_TOTAL_MEMBER: raise ValueError(f"daily contract must select {DANIEL_MOSKOWITZ_DAILY_TOTAL_MEMBER!r}") if self.monthly.member_name != DANIEL_MOSKOWITZ_MONTHLY_TOTAL_MEMBER: raise ValueError( f"monthly contract must select {DANIEL_MOSKOWITZ_MONTHLY_TOTAL_MEMBER!r}" ) _optional_digest(self.expected_archive_sha256, label="expected_archive_sha256") tolerance = self.daily_monthly_alignment_atol if tolerance is not None and ( isinstance(tolerance, bool) or not isinstance(tolerance, (int, float)) or not math.isfinite(float(tolerance)) or float(tolerance) < 0.0 ): raise ValueError("daily_monthly_alignment_atol must be finite, nonnegative, or None")
[docs] @dataclass(frozen=True) class DanielMoskowitzMomentumData: """Canonical paired momentum-decile frames returned by the path-only loader.""" daily: pd.DataFrame monthly: pd.DataFrame
@dataclass(frozen=True) class _ArchiveContents: archive_sha256: str member_manifest_sha256: str selected: dict[str, bytes] member_sha256: dict[str, str] member_bytes: dict[str, int] def _local_bytes(path: str | Path) -> bytes: source = Path(path).expanduser() if ".." in source.parts: raise ValueError("Daniel--Moskowitz archive path must not contain '..' components") absolute = Path(os.path.abspath(source)) if tuple(suffix.casefold() for suffix in absolute.suffixes[-2:]) != (".tar", ".gz"): raise ValueError("Daniel--Moskowitz input must be a .tar.gz 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 Daniel--Moskowitz 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( "Daniel--Moskowitz 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( "Daniel--Moskowitz 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( "Daniel--Moskowitz 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("Daniel--Moskowitz archive path is not a regular file") if before.st_size <= 0: raise ValueError("Daniel--Moskowitz archive is empty") if before.st_size > _MAX_ARCHIVE_BYTES: raise ValueError("Daniel--Moskowitz 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("Daniel--Moskowitz archive path is not readable") from exc finally: os.close(file_fd) if len(content) > _MAX_ARCHIVE_BYTES: raise ValueError("Daniel--Moskowitz 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("Daniel--Moskowitz archive changed while it was being read") return bytes(content) def _safe_member_name(name: str) -> str: if not name or "\\" in name or "\0" in name: raise ValueError("Daniel--Moskowitz archive contains an unsafe member name") path = PurePosixPath(name) if path.is_absolute() or any(part in {"", ".", ".."} for part in path.parts): raise ValueError("Daniel--Moskowitz archive contains an unsafe member path") canonical = path.as_posix() if canonical != name: raise ValueError("Daniel--Moskowitz archive contains a non-canonical member path") if len(path.parts) != 1: raise ValueError("Daniel--Moskowitz archive members must be root-level files") return canonical def _tar_octal_size(field: bytes) -> int: stripped = field.strip(b" \0") if not stripped: return 0 if any(value < ord("0") or value > ord("7") for value in stripped): raise ValueError("Daniel--Moskowitz archive uses an unsupported tar size encoding") return int(stripped, 8) def _gzip_read_exact( stream: gzip.GzipFile, size: int, *, stream_bytes: int, collect: bool, ) -> tuple[bytes, int]: chunks: list[bytes] = [] remaining = size while remaining: chunk = stream.read(min(_READ_CHUNK_BYTES, remaining)) if not chunk: raise ValueError("Daniel--Moskowitz tar stream is truncated") remaining -= len(chunk) stream_bytes += len(chunk) if stream_bytes > _MAX_UNCOMPRESSED_BYTES + _MAX_TAR_OVERHEAD_BYTES: raise ValueError("Daniel--Moskowitz tar stream exceeds the safety limit") if collect: chunks.append(chunk) return b"".join(chunks), stream_bytes def _validate_raw_tar_envelope(content: bytes) -> int: """Reject hidden metadata headers and nonzero payload after the tar end marker.""" raw = io.BytesIO(content) stream_bytes = 0 member_names: list[str] = [] try: with gzip.GzipFile(fileobj=raw, mode="rb") as stream: while True: header, stream_bytes = _gzip_read_exact( stream, tarfile.BLOCKSIZE, stream_bytes=stream_bytes, collect=True, ) if not any(header): second, stream_bytes = _gzip_read_exact( stream, tarfile.BLOCKSIZE, stream_bytes=stream_bytes, collect=True, ) if any(second): raise ValueError("Daniel--Moskowitz tar stream has an invalid end marker") trailing = 0 while True: chunk = stream.read(_READ_CHUNK_BYTES) if not chunk: break stream_bytes += len(chunk) trailing += len(chunk) if trailing > _MAX_TAR_TRAILING_ZERO_BYTES: raise ValueError( "Daniel--Moskowitz tar stream has excessive trailing padding" ) if any(chunk): raise ValueError( "Daniel--Moskowitz tar stream contains a trailing payload" ) break type_flag = header[156:157] if type_flag == tarfile.SYMTYPE: raise ValueError("Daniel--Moskowitz archive must not contain symbolic links") if type_flag == tarfile.LNKTYPE: raise ValueError("Daniel--Moskowitz archive must not contain hard links") metadata_types = { tarfile.GNUTYPE_LONGLINK, tarfile.GNUTYPE_LONGNAME, tarfile.SOLARIS_XHDTYPE, tarfile.XGLTYPE, tarfile.XHDTYPE, } if type_flag in metadata_types: raise ValueError( "Daniel--Moskowitz archive must not contain tar metadata extension headers" ) if type_flag not in {tarfile.REGTYPE, tarfile.AREGTYPE}: raise ValueError("Daniel--Moskowitz archive must contain only regular files") raw_name = header[:100].split(b"\0", 1)[0] raw_prefix = header[345:500].split(b"\0", 1)[0] try: name = raw_name.decode("ascii") if raw_prefix: name = f"{raw_prefix.decode('ascii')}/{name}" except UnicodeError as exc: raise ValueError( "Daniel--Moskowitz archive member names must be ASCII" ) from exc name = _safe_member_name(name) if name not in _ALLOWED_MEMBERS: raise ValueError( f"Daniel--Moskowitz archive contains unexpected member {name!r}" ) member_names.append(name) if len(member_names) > _MAX_MEMBERS: raise ValueError( "Daniel--Moskowitz archive has an invalid number of archive members" ) member_size = _tar_octal_size(header[124:136]) if member_size <= 0 or member_size > _MAX_MEMBER_BYTES: raise ValueError( "Daniel--Moskowitz archive member has an invalid or unsafe size" ) _, stream_bytes = _gzip_read_exact( stream, member_size, stream_bytes=stream_bytes, collect=False, ) padding_size = (-member_size) % tarfile.BLOCKSIZE padding, stream_bytes = _gzip_read_exact( stream, padding_size, stream_bytes=stream_bytes, collect=True, ) if any(padding): raise ValueError("Daniel--Moskowitz archive member has nonzero tar padding") except (gzip.BadGzipFile, EOFError, OSError) as exc: raise ValueError("Daniel--Moskowitz input is not a valid gzip tar archive") from exc folded = [name.casefold() for name in member_names] if len(folded) != len(set(folded)): raise ValueError("Daniel--Moskowitz archive has duplicate member names") if set(folded) != {name.casefold() for name in _ALLOWED_MEMBERS}: raise ValueError("Daniel--Moskowitz archive must contain exactly the documented member set") if raw.tell() != len(content): raise ValueError("Daniel--Moskowitz gzip envelope contains unread trailing bytes") return stream_bytes def _read_archive(content: bytes, *, selected_names: frozenset[str]) -> _ArchiveContents: archive_digest = hashlib.sha256(content).hexdigest() tar_stream_bytes = _validate_raw_tar_envelope(content) if tar_stream_bytes / len(content) > _MAX_COMPRESSION_RATIO: raise ValueError("Daniel--Moskowitz archive has an unsafe compression ratio") selected: dict[str, bytes] = {} digests: dict[str, str] = {} sizes: dict[str, int] = {} try: with tarfile.open(fileobj=io.BytesIO(content), mode="r:gz") as archive: members = archive.getmembers() if not members or len(members) > _MAX_MEMBERS: raise ValueError( "Daniel--Moskowitz archive has an invalid number of archive members" ) seen: set[str] = set() total = 0 validated: list[tuple[tarfile.TarInfo, str]] = [] for member in members: name = _safe_member_name(member.name) folded = name.casefold() if folded in seen: raise ValueError("Daniel--Moskowitz archive has duplicate member names") seen.add(folded) if member.issym(): raise ValueError("Daniel--Moskowitz archive must not contain symbolic links") if member.islnk(): raise ValueError("Daniel--Moskowitz archive must not contain hard links") if not member.isreg(): raise ValueError("Daniel--Moskowitz archive must contain only regular files") if name not in _ALLOWED_MEMBERS: raise ValueError( f"Daniel--Moskowitz archive contains unexpected member {name!r}" ) if member.size <= 0 or member.size > _MAX_MEMBER_BYTES: raise ValueError( "Daniel--Moskowitz archive member has an invalid or unsafe size" ) total += member.size if total > _MAX_UNCOMPRESSED_BYTES: raise ValueError( "Daniel--Moskowitz archive exceeds the uncompressed-size safety limit" ) validated.append((member, name)) if seen != {name.casefold() for name in _ALLOWED_MEMBERS}: raise ValueError( "Daniel--Moskowitz archive must contain exactly the documented member set" ) for member, name in validated: extracted = archive.extractfile(member) if extracted is None: raise ValueError(f"Daniel--Moskowitz archive member {name!r} could not be read") digest = hashlib.sha256() chunks: list[bytes] | None = [] if name in selected_names else None count = 0 while True: chunk = extracted.read(_READ_CHUNK_BYTES) if not chunk: break count += len(chunk) if count > member.size or count > _MAX_MEMBER_BYTES: raise ValueError( f"Daniel--Moskowitz archive member {name!r} exceeds its declared size" ) digest.update(chunk) if chunks is not None: chunks.append(chunk) if count != member.size: raise ValueError(f"Daniel--Moskowitz archive member {name!r} is truncated") digests[name] = digest.hexdigest() sizes[name] = count if chunks is not None: selected[name] = b"".join(chunks) except (tarfile.TarError, EOFError, OSError) as exc: raise ValueError("Daniel--Moskowitz input is not a valid gzip tar archive") from exc manifest_payload = [ {"name": name, "sha256": digests[name], "size": sizes[name]} for name in sorted(digests) ] manifest = hashlib.sha256( json.dumps(manifest_payload, sort_keys=True, separators=(",", ":")).encode("ascii") ).hexdigest() return _ArchiveContents( archive_sha256=archive_digest, member_manifest_sha256=manifest, selected=selected, member_sha256=digests, member_bytes=sizes, ) def _bounded_physical_lines(text: str, *, member_name: str) -> Iterator[str]: for line_number, line in enumerate(io.StringIO(text), start=1): if line_number > _MAX_TEXT_PHYSICAL_LINES: raise ValueError( f"Daniel--Moskowitz member {member_name!r} exceeds the physical-line limit" ) if len(line) > _MAX_TEXT_RECORD_CHARS: raise ValueError( f"Daniel--Moskowitz member {member_name!r} exceeds the line-length limit" ) if any(character.isspace() and character not in _ASCII_WHITESPACE for character in line): raise ValueError( f"Daniel--Moskowitz member {member_name!r} contains non-ASCII whitespace" ) yield line def _bounded_records(rows: Iterable[list[str]], *, member_name: str) -> list[list[str]]: records: list[list[str]] = [] field_count = 0 for raw_row in rows: row = [cell.strip(_ASCII_WHITESPACE) for cell in raw_row] if not any(row): continue if len(records) >= _MAX_TEXT_RECORDS: raise ValueError( f"Daniel--Moskowitz member {member_name!r} exceeds the text-record limit" ) if len(row) > _MAX_TEXT_COLUMNS: raise ValueError( f"Daniel--Moskowitz member {member_name!r} exceeds the column-count limit" ) if sum(len(cell) for cell in row) > _MAX_TEXT_RECORD_CHARS: raise ValueError( f"Daniel--Moskowitz member {member_name!r} exceeds the record-length limit" ) field_count += len(row) if field_count > _MAX_TEXT_FIELDS: raise ValueError( f"Daniel--Moskowitz member {member_name!r} exceeds the aggregate-field limit" ) records.append(row) return records def _records(content: bytes, *, contract: DanielMoskowitzMemberContract) -> list[list[str]]: try: text = content.decode(contract.encoding, errors="strict") except UnicodeError as exc: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} is not valid {contract.encoding}" ) from exc if "\0" in text: raise ValueError(f"Daniel--Moskowitz member {contract.member_name!r} contains a NUL byte") lines = _bounded_physical_lines(text, member_name=contract.member_name) if contract.delimiter == "whitespace": return _bounded_records( ( [] if not (stripped := line.strip(_ASCII_WHITESPACE)) else _ASCII_WHITESPACE_PATTERN.split(stripped) for line in lines ), member_name=contract.member_name, ) try: reader = csv.reader( lines, delimiter=contract.delimiter, strict=True, ) return _bounded_records(reader, member_name=contract.member_name) except csv.Error as exc: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has malformed delimited text" ) from exc def _resolve_column( reference: DanielMoskowitzColumn, *, headers: tuple[str, ...] | None, width: int, member_name: str, ) -> int: if isinstance(reference, str): if headers is None: # guarded by contract validation; retained as a fail-closed invariant raise ValueError("string column references require a header row") try: return headers.index(reference) except ValueError: raise ValueError( f"Daniel--Moskowitz member {member_name!r} has no column {reference!r}" ) from None if reference >= width: raise ValueError( f"Daniel--Moskowitz member {member_name!r} has no column position {reference}" ) return reference def _strict_source_date(value: str, *, date_format: str, member_name: str, row: int) -> datetime: try: parsed = datetime.strptime(value, date_format) except ValueError as exc: raise ValueError( f"Daniel--Moskowitz member {member_name!r} has an invalid date in data row {row}" ) from exc if parsed.strftime(date_format) != value: raise ValueError( f"Daniel--Moskowitz member {member_name!r} has a non-canonical date in data row {row}" ) if parsed != datetime(parsed.year, parsed.month, parsed.day): raise ValueError(f"Daniel--Moskowitz member {member_name!r} dates must not contain a time") return parsed def _expected_date( value: str, *, frequency: Literal["daily", "monthly"], label: str ) -> pd.Timestamp: try: stamp = pd.Timestamp(value) except (TypeError, ValueError) as exc: raise ValueError(f"{label} must be a valid timezone-naive date") from exc if pd.isna(stamp) or stamp.tz is not None: raise ValueError(f"{label} must be a valid timezone-naive date") normalized = stamp.normalize() if stamp != normalized: raise ValueError(f"{label} must not contain a time") if frequency == "monthly": return normalized.to_period("M").to_timestamp("M") return normalized def _parse_member( content: bytes, *, contract: DanielMoskowitzMemberContract, frequency: Literal["daily", "monthly"], ) -> pd.DataFrame: records = _records(content, contract=contract) if not records: raise ValueError(f"Daniel--Moskowitz member {contract.member_name!r} is empty") header_index = contract.header_row headers: tuple[str, ...] | None = None if header_index is not None: if header_index >= len(records): raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has no declared header row" ) headers = tuple(cell.strip(_ASCII_WHITESPACE) for cell in records[header_index]) if any(not header for header in headers) or len( {header.casefold() for header in headers} ) != len(headers): raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has empty or duplicate headers" ) start = contract.data_start_row if start is None: start = 0 if header_index is None else header_index + 1 end = len(records) if contract.data_end_row is None else contract.data_end_row if start >= len(records) or end > len(records) or end <= start: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has an invalid data-row envelope" ) data = records[start:end] width = len(headers) if headers is not None else len(data[0]) if width == 0 or any(len(row) != width for row in data): raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has inconsistent row widths" ) date_position = _resolve_column( contract.date_column, headers=headers, width=width, member_name=contract.member_name, ) portfolio_positions = tuple( _resolve_column( column, headers=headers, width=width, member_name=contract.member_name, ) for column in contract.portfolio_columns ) if len(set(portfolio_positions)) != 10 or date_position in portfolio_positions: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} resolves duplicate source columns" ) dates: list[pd.Timestamp] = [] values = np.empty((len(data), 10), dtype=np.float64) divisor = 100.0 if contract.source_unit == "percent" else 1.0 for row_number, row in enumerate(data, start=1): parsed = _strict_source_date( row[date_position], date_format=contract.date_format, member_name=contract.member_name, row=row_number, ) stamp = pd.Timestamp(parsed) if frequency == "monthly": stamp = stamp.to_period("M").to_timestamp("M") dates.append(stamp) for column_number, position in enumerate(portfolio_positions): raw = row[position] try: value = float(raw) except ValueError as exc: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has a non-numeric " f"return in data row {row_number}" ) from exc value /= divisor if not math.isfinite(value): raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has a non-finite return" ) if value < -1.0 or value > _MAX_DECIMAL_LONG_ONLY_RETURN: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} violates the declared " "simple-return unit envelope" ) values[row_number - 1, column_number] = value date_index = pd.DatetimeIndex(dates, name="date") if date_index.has_duplicates: raise ValueError(f"Daniel--Moskowitz member {contract.member_name!r} has duplicate dates") if not date_index.is_monotonic_increasing: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} dates are not strictly increasing" ) if frequency == "monthly": periods = date_index.to_period("M") expected_periods = pd.period_range(periods[0], periods[-1], freq="M") if not periods.equals(expected_periods): raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} has a monthly calendar gap" ) if contract.expected_rows is not None and len(date_index) != contract.expected_rows: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} violates expected_rows" ) if contract.expected_start_date is not None: expected_start = _expected_date( contract.expected_start_date, frequency=frequency, label="expected_start_date", ) if date_index[0] != expected_start: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} violates expected_start_date" ) if contract.expected_end_date is not None: expected_end = _expected_date( contract.expected_end_date, frequency=frequency, label="expected_end_date", ) if date_index[-1] != expected_end: raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} violates expected_end_date" ) if contract.expected_dates is not None: expected_dates = pd.DatetimeIndex( [ _expected_date(value, frequency=frequency, label="expected_dates") for value in contract.expected_dates ] ) if not date_index.equals(expected_dates): raise ValueError( f"Daniel--Moskowitz member {contract.member_name!r} violates expected_dates" ) columns = [f"decile_{index}" for index in range(1, 11)] frame = pd.DataFrame(values, columns=columns) frame.insert(0, "date", date_index.to_numpy()) frame["wml"] = frame["decile_10"] - frame["decile_1"] return frame def _validate_alignment( daily: pd.DataFrame, monthly: pd.DataFrame, *, tolerance: float | None, ) -> None: daily_months = pd.PeriodIndex(daily["date"], freq="M").unique() monthly_months = pd.PeriodIndex(monthly["date"], freq="M") if not daily_months.equals(monthly_months): raise ValueError( "Daniel--Moskowitz daily and monthly members cover different calendar months" ) if tolerance is None: return columns = [f"decile_{index}" for index in range(1, 11)] grouped = daily.assign(_month=pd.PeriodIndex(daily["date"], freq="M")).groupby( "_month", sort=False ) compounded = grouped[columns].agg(lambda series: float(np.prod(1.0 + series) - 1.0)) observed = monthly.set_index(pd.PeriodIndex(monthly["date"], freq="M"))[columns] if not np.allclose( compounded.to_numpy(dtype=np.float64), observed.to_numpy(dtype=np.float64), rtol=0.0, atol=float(tolerance), ): raise ValueError( "Daniel--Moskowitz daily returns do not compound to the monthly member within " "daily_monthly_alignment_atol" ) def _contract_digest(contract: DanielMoskowitzArchiveContract) -> str: payload = json.dumps(asdict(contract), sort_keys=True, separators=(",", ":")) return hashlib.sha256(payload.encode("utf-8")).hexdigest() def _provenance( *, frequency: Literal["daily", "monthly"], contract: DanielMoskowitzArchiveContract, member_contract: DanielMoskowitzMemberContract, paired_contract: DanielMoskowitzMemberContract, contents: _ArchiveContents, frame: pd.DataFrame, contract_sha256: str, ) -> dict[str, str]: archive_pinned = contract.expected_archive_sha256 is not None member_pinned = member_contract.expected_sha256 is not None identity = ( "expected_archive_and_member_sha256_verified" if archive_pinned and member_pinned else "actual_sha256_recorded_with_incomplete_caller_pins" ) data_vintage = ( f"daniel-moskowitz-momentum@sha256:{contents.archive_sha256}" f"/contract:{contract_sha256}/recipe:{_RECIPE}" ) return { "backend": "daniel-moskowitz-local-archive", "dataset": "Daniel-Moskowitz Momentum Portfolios", "source_page": DANIEL_MOSKOWITZ_DATA_PAGE, "documentation_url": DANIEL_MOSKOWITZ_DOCUMENTATION_URL, "expected_archive_filename": DANIEL_MOSKOWITZ_ARCHIVE_FILENAME, "archive_sha256": contents.archive_sha256, "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_contract.member_name, "member_sha256": contents.member_sha256[member_contract.member_name], "member_bytes": str(contents.member_bytes[member_contract.member_name]), "paired_member_name": paired_contract.member_name, "paired_member_sha256": contents.member_sha256[paired_contract.member_name], "frequency": frequency, "source_unit": member_contract.source_unit, "unit_contract": f"source_{member_contract.source_unit}_to_decimal_simple_returns", "full_start_date": pd.Timestamp(frame["date"].iloc[0]).date().isoformat(), "full_end_date": pd.Timestamp(frame["date"].iloc[-1]).date().isoformat(), "full_rows": str(len(frame)), "recipe": _RECIPE, "contract_sha256": contract_sha256, "identity_validation": identity, "paper_exact_claim": "not_established_by_loader", "market_state_inputs": "absent-from-author-archive", "required_market_state_inputs": ( "monthly_market_total_return,daily_market_excess_return,risk_free" ), "vintage_semantics": "static_caller_obtained_archive_non_pit", "redistributable": "false", "data_vintage": data_vintage, } def load_archive( path: str | Path, *, contract: DanielMoskowitzArchiveContract, ) -> tuple[DanielMoskowitzMomentumData, dict[str, str], dict[str, str]]: """Parse the paired total-return members without extracting or persisting archive bytes.""" if not isinstance(contract, DanielMoskowitzArchiveContract): raise ValueError("contract must be a DanielMoskowitzArchiveContract") content = _local_bytes(path) contents = _read_archive( content, selected_names=frozenset({contract.daily.member_name, contract.monthly.member_name}), ) if ( contract.expected_archive_sha256 is not None and contents.archive_sha256 != contract.expected_archive_sha256.casefold() ): raise ValueError("Daniel--Moskowitz archive does not match expected_archive_sha256") for member in (contract.daily, contract.monthly): if ( member.expected_sha256 is not None and contents.member_sha256[member.member_name] != member.expected_sha256.casefold() ): raise ValueError( f"Daniel--Moskowitz member {member.member_name!r} does not match expected_sha256" ) daily = _parse_member( contents.selected[contract.daily.member_name], contract=contract.daily, frequency="daily", ) monthly = _parse_member( contents.selected[contract.monthly.member_name], contract=contract.monthly, frequency="monthly", ) _validate_alignment( daily, monthly, tolerance=contract.daily_monthly_alignment_atol, ) contract_sha256 = _contract_digest(contract) daily_provenance = _provenance( frequency="daily", contract=contract, member_contract=contract.daily, paired_contract=contract.monthly, contents=contents, frame=daily, contract_sha256=contract_sha256, ) monthly_provenance = _provenance( frequency="monthly", contract=contract, member_contract=contract.monthly, paired_contract=contract.daily, contents=contents, frame=monthly, contract_sha256=contract_sha256, ) parsed = DanielMoskowitzMomentumData(daily=daily, monthly=monthly) return parsed, daily_provenance, monthly_provenance