From 76d4f4edb7cf04b99b43df51f37bddf279146d46 Mon Sep 17 00:00:00 2001 From: HJ <16863475+hjcud@users.noreply.github.com> Date: Sun, 9 Aug 2026 18:35:20 +0900 Subject: [PATCH] fix: consume immutable legacy market manifests --- src/backtest_engine/execution_policy.py | 31 ++- src/backtest_engine/legacy_market_data.py | 262 ++++++++++++++++++ src/backtest_engine/market_data.py | 52 +++- src/backtest_engine/production.py | 36 ++- src/backtest_engine/wiring.py | 22 +- .../int03-development-market-manifest.json | 79 ++++++ tests/test_execution_policy.py | 20 +- tests/test_legacy_market_data.py | 151 ++++++++++ tests/test_production.py | 81 +++++- 9 files changed, 692 insertions(+), 42 deletions(-) create mode 100644 src/backtest_engine/legacy_market_data.py create mode 100644 tests/fixtures/int03-development-market-manifest.json create mode 100644 tests/test_legacy_market_data.py diff --git a/src/backtest_engine/execution_policy.py b/src/backtest_engine/execution_policy.py index d3fc374..ebcc051 100644 --- a/src/backtest_engine/execution_policy.py +++ b/src/backtest_engine/execution_policy.py @@ -12,8 +12,9 @@ The ET calendar quarter of strategy releases this policy governs. It is the catalog selection key. ``period_start`` / ``period_end`` - The pinned evaluation window. It must span whole ET calendar quarters; a - policy cannot present a 60-minute window as a quarterly official backtest. + The pinned evaluation window. It is independent of ``release_quarter``: + Development may publish a shorter immutable evaluation period when that is + the complete locked dataset available for the release. Order horizons -------------- @@ -33,7 +34,7 @@ from datetime import datetime, timedelta, timezone from decimal import Decimal from types import MappingProxyType -from zoneinfo import ZoneInfo +from zoneinfo import ZoneInfo, ZoneInfoNotFoundError from .money import PRECISION_RULES_VERSION @@ -125,14 +126,17 @@ def __post_init__(self) -> None: raise ValueError("execution policy period must use UTC") if self.period_start >= self.period_end: raise ValueError("execution policy period must be increasing") - - start_index = _quarter_index(self.period_start) - end_index = _quarter_index(self.period_end) - if start_index is None or end_index is None: - raise ValueError( - "execution policy period must span whole ET calendar quarter " - "boundaries" - ) + try: + policy_zone = ZoneInfo(self.timezone) + except (ValueError, ZoneInfoNotFoundError) as exc: + raise ValueError(f"execution policy timezone is invalid: {self.timezone!r}") from exc + for field_name in ("period_start", "period_end"): + local = getattr(self, field_name).astimezone(policy_zone) + if (local.hour, local.minute, local.second, local.microsecond) != (0, 0, 0, 0): + raise ValueError( + "execution policy period must use local calendar-day boundaries: " + f"{field_name}={local.isoformat()}" + ) if self.fee_rate < 0: raise ValueError("execution policy rates must not be negative") @@ -171,10 +175,11 @@ def __post_init__(self) -> None: @property def quarter_count(self) -> int: - """Number of whole ET calendar quarters covered by the period.""" + """Number of whole ET quarters, for quarter-aligned policy fixtures.""" start_index = _quarter_index(self.period_start) end_index = _quarter_index(self.period_end) - assert start_index is not None and end_index is not None + if start_index is None or end_index is None: + raise ValueError("execution policy period is not aligned to ET quarter boundaries") return end_index - start_index @property diff --git a/src/backtest_engine/legacy_market_data.py b/src/backtest_engine/legacy_market_data.py new file mode 100644 index 0000000..ba7fdf3 --- /dev/null +++ b/src/backtest_engine/legacy_market_data.py @@ -0,0 +1,262 @@ +"""Fail-closed adapter for immutable ``market-loader/1.0.0`` manifests. + +The Development catalog predates the canonical ``market-data.v1`` object +shape. Its rows are still reproducible, but their dataset hash covers the +loader's logical publication material rather than the current canonical object +metadata. This module recognizes only that exact legacy shape and recomputes +the original producer hash. It is deliberately separate from the canonical +validator so accepting the old Development fixture cannot widen the current +contract. +""" + +from __future__ import annotations + +import hashlib +import json +import re +import uuid +from collections.abc import Mapping, Sequence +from datetime import date, datetime +from typing import Any, cast +from zoneinfo import ZoneInfo + + +__all__ = [ + "LEGACY_MARKET_SCHEMA_ID", + "LegacyMarketDataError", + "is_legacy_market_loader_manifest", + "legacy_dataset_hash", + "legacy_period_matches", + "validate_legacy_market_loader_manifest", +] + + +LEGACY_MARKET_SCHEMA_ID = "market-bars/1" +LEGACY_PROCESSING_VERSION = "market-loader/1.0.0" +_SHA256 = re.compile(r"^[0-9a-f]{64}$") +_KEY = re.compile( + r"^historical/provider=alpaca/feed=sip/adjustment=(?Pall)/" + r"session=regular/resolution=(?P30m|1h|4h|1d)/" + r"revision=(?P[0-9]{8})/year=(?P[0-9]{4})/" + r"shard=(?P[0-9]{2})-of-(?P[0-9]{2})/" + r"manifest_id=(?P[0-9a-f-]{36})/part-(?P[0-9]{5})\.parquet$" +) + + +class LegacyMarketDataError(ValueError): + """The legacy catalog evidence cannot identify one immutable publication.""" + + +def _mapping(value: object, label: str) -> Mapping[str, Any]: + if not isinstance(value, Mapping): + raise LegacyMarketDataError(f"{label} must be an object") + return value + + +def _text(value: object, label: str) -> str: + if not isinstance(value, str) or not value: + raise LegacyMarketDataError(f"{label} must be a non-empty string") + return value + + +def _uuid(value: object, label: str) -> str: + text = _text(value, label) + try: + parsed = uuid.UUID(text) + except ValueError as exc: + raise LegacyMarketDataError(f"{label} must be a UUID") from exc + if str(parsed) != text: + raise LegacyMarketDataError(f"{label} must use canonical lowercase UUID text") + return text + + +def _period_date(value: object, label: str) -> date: + text = _text(value, label) + if not text.endswith("T00:00:00Z"): + raise LegacyMarketDataError(f"{label} must be a legacy UTC date boundary") + try: + return date.fromisoformat(text[:10]) + except ValueError as exc: + raise LegacyMarketDataError(f"{label} must contain an ISO date") from exc + + +def is_legacy_market_loader_manifest(manifest: Mapping[str, Any]) -> bool: + """Return whether the document explicitly selects the isolated legacy path.""" + return manifest.get("schema_id") == LEGACY_MARKET_SCHEMA_ID + + +def _legacy_objects(manifest: Mapping[str, Any]) -> list[dict[str, object]]: + manifest_id = _uuid(manifest.get("manifest_id"), "manifest_id") + revision = manifest.get("revision") + if not isinstance(revision, int) or isinstance(revision, bool) or revision < 1: + raise LegacyMarketDataError("revision must be a positive integer") + resolution = _text(manifest.get("resolution"), "resolution") + period_start = _period_date(manifest.get("period_start"), "period_start") + period_end = _period_date(manifest.get("period_end"), "period_end") + if period_start >= period_end: + raise LegacyMarketDataError("legacy manifest period must be increasing") + raw_objects = manifest.get("objects") + if not isinstance(raw_objects, Sequence) or isinstance(raw_objects, (str, bytes)): + raise LegacyMarketDataError("objects must be a non-empty array") + if not raw_objects: + raise LegacyMarketDataError("objects must be a non-empty array") + + hashed: list[dict[str, object]] = [] + seen: set[tuple[int, int]] = set() + shard_counts: set[int] = set() + for index, raw in enumerate(raw_objects): + item = _mapping(raw, f"objects[{index}]") + _uuid(item.get("storage_object_id"), f"objects[{index}].storage_object_id") + key = _text(item.get("object_key"), f"objects[{index}].object_key") + match = _KEY.fullmatch(key) + if match is None: + raise LegacyMarketDataError(f"objects[{index}].object_key is not a legacy loader key") + if match["manifest_id"] != manifest_id: + raise LegacyMarketDataError(f"objects[{index}].object_key binds another manifest") + if int(match["revision"]) != revision: + raise LegacyMarketDataError(f"objects[{index}].object_key revision does not match") + if match["resolution"] != resolution: + raise LegacyMarketDataError(f"objects[{index}].object_key resolution does not match") + if int(match["year"]) != period_start.year: + raise LegacyMarketDataError(f"objects[{index}].object_key year does not match") + + shard = int(match["shard"]) + shard_count = int(match["shard_count"]) + part = int(match["part"]) + if shard_count < 1 or shard >= shard_count or part < 1: + raise LegacyMarketDataError(f"objects[{index}].object_key shard/part is invalid") + if (shard, part) in seen: + raise LegacyMarketDataError("legacy manifest contains a duplicate shard/part") + seen.add((shard, part)) + shard_counts.add(shard_count) + if item.get("shard_key") != f"s{shard:02d}-of-{shard_count}": + raise LegacyMarketDataError(f"objects[{index}].shard_key does not match object_key") + if item.get("part_number") != part: + raise LegacyMarketDataError(f"objects[{index}].part_number does not match object_key") + + content_hash = _text(item.get("content_hash"), f"objects[{index}].content_hash") + if _SHA256.fullmatch(content_hash) is None: + raise LegacyMarketDataError(f"objects[{index}].content_hash must be lowercase SHA-256") + row_count = item.get("row_count") + if not isinstance(row_count, int) or isinstance(row_count, bool) or row_count < 0: + raise LegacyMarketDataError(f"objects[{index}].row_count must be non-negative") + if item.get("object_kind") != "MARKET_BARS": + raise LegacyMarketDataError(f"objects[{index}].object_kind must be MARKET_BARS") + if item.get("partition_granularity") != "YEAR": + raise LegacyMarketDataError(f"objects[{index}].partition_granularity must be YEAR") + if item.get("schema_version") != LEGACY_MARKET_SCHEMA_ID: + raise LegacyMarketDataError(f"objects[{index}].schema_version does not match") + if item.get("partition_start") != period_start.isoformat(): + raise LegacyMarketDataError(f"objects[{index}].partition_start does not match") + if item.get("partition_end") != period_end.isoformat(): + raise LegacyMarketDataError(f"objects[{index}].partition_end does not match") + if _period_date(item.get("period_start"), f"objects[{index}].period_start") != period_start: + raise LegacyMarketDataError(f"objects[{index}].period_start does not match") + if _period_date(item.get("period_end"), f"objects[{index}].period_end") != period_end: + raise LegacyMarketDataError(f"objects[{index}].period_end does not match") + + storage_fields = ( + item.get("storage_provider"), + item.get("bucket_name"), + item.get("provider_version_id"), + ) + if any(value is not None for value in storage_fields): + if item.get("storage_provider") != "S3": + raise LegacyMarketDataError(f"objects[{index}].storage_provider must be S3") + _text(item.get("bucket_name"), f"objects[{index}].bucket_name") + _text(item.get("provider_version_id"), f"objects[{index}].provider_version_id") + + hashed.append( + { + "content_sha256": content_hash, + "row_count": row_count, + "period_start": period_start.isoformat(), + "period_end": period_end.isoformat(), + "shard": shard, + "part": part, + } + ) + + if len(shard_counts) != 1: + raise LegacyMarketDataError("legacy manifest objects disagree on shard count") + shard_count = shard_counts.pop() + if seen != {(shard, 1) for shard in range(shard_count)}: + raise LegacyMarketDataError("legacy manifest must contain exactly one part for every shard") + return sorted( + hashed, + key=lambda item: (cast(int, item["shard"]), cast(int, item["part"])), + ) + + +def legacy_dataset_hash(manifest: Mapping[str, Any]) -> str: + """Recompute the exact ``market-loader/1.0.0`` publication digest.""" + payload = { + "provider": "ALPACA", + "feed": _text(manifest.get("feed_code"), "feed_code"), + "adjustment": "all", + "session": "XNYS_REGULAR", + "resolution": _text(manifest.get("resolution"), "resolution"), + "period_start": _period_date(manifest.get("period_start"), "period_start").isoformat(), + "period_end": _period_date(manifest.get("period_end"), "period_end").isoformat(), + "revision": manifest.get("revision"), + "schema_version": manifest.get("schema_id"), + "processing_version": LEGACY_PROCESSING_VERSION, + "objects": _legacy_objects(manifest), + } + encoded = json.dumps(payload, sort_keys=True, separators=(",", ":"), ensure_ascii=False) + return hashlib.sha256(encoded.encode("utf-8")).hexdigest() + + +def validate_legacy_market_loader_manifest(manifest: Mapping[str, Any]) -> None: + """Validate the isolated legacy shape and its original producer digest.""" + if manifest.get("contract_id") != "com06.dataset-manifest": + raise LegacyMarketDataError("contract_id must be com06.dataset-manifest") + if manifest.get("schema_version") != 1: + raise LegacyMarketDataError("schema_version must be 1") + manifest_id = _uuid(manifest.get("manifest_id"), "manifest_id") + if _uuid(manifest.get("dataset_id"), "dataset_id") != manifest_id: + raise LegacyMarketDataError("legacy dataset_id must equal manifest_id") + if manifest.get("schema_id") != LEGACY_MARKET_SCHEMA_ID: + raise LegacyMarketDataError(f"schema_id must be {LEGACY_MARKET_SCHEMA_ID}") + if manifest.get("provider_code") != "ALPACA": + raise LegacyMarketDataError("provider_code must be ALPACA") + resolution = _text(manifest.get("resolution"), "resolution") + if manifest.get("feed_code") != f"ALPACA_SIP_ALL_{resolution.upper()}": + raise LegacyMarketDataError("feed_code does not match the adjusted SIP resolution") + if manifest.get("data_layer") != "ADJUSTED": + raise LegacyMarketDataError("data_layer must be ADJUSTED") + if manifest.get("status") != "AVAILABLE": + raise LegacyMarketDataError("legacy manifest must be AVAILABLE") + available_at = _text(manifest.get("available_at"), "available_at") + try: + datetime.fromisoformat(available_at.replace("Z", "+00:00")) + except ValueError as exc: + raise LegacyMarketDataError("available_at must be ISO-8601") from exc + declared = _text(manifest.get("dataset_hash"), "dataset_hash") + if _SHA256.fullmatch(declared) is None: + raise LegacyMarketDataError("dataset_hash must be lowercase SHA-256") + computed = legacy_dataset_hash(manifest) + if declared != computed: + raise LegacyMarketDataError( + "legacy dataset_hash does not match market-loader/1.0.0 publication material: " + f"declared {declared}, computed {computed}" + ) + + +def legacy_period_matches( + manifest: Mapping[str, Any], + period_start: datetime, + period_end: datetime, + timezone_name: str, +) -> bool: + """Compare loader date labels with policy-local midnight boundaries.""" + zone = ZoneInfo(timezone_name) + local_start = period_start.astimezone(zone) + local_end = period_end.astimezone(zone) + midnight = (0, 0, 0, 0) + return ( + (local_start.hour, local_start.minute, local_start.second, local_start.microsecond) == midnight + and (local_end.hour, local_end.minute, local_end.second, local_end.microsecond) == midnight + and _period_date(manifest.get("period_start"), "period_start") == local_start.date() + and _period_date(manifest.get("period_end"), "period_end") == local_end.date() + ) diff --git a/src/backtest_engine/market_data.py b/src/backtest_engine/market_data.py index ceec8f1..ecbe9a6 100644 --- a/src/backtest_engine/market_data.py +++ b/src/backtest_engine/market_data.py @@ -30,6 +30,12 @@ from .contracts import ContractValidationError, validate_dataset_manifest from .execution_policy import ExecutionPolicy +from .legacy_market_data import ( + LEGACY_MARKET_SCHEMA_ID, + is_legacy_market_loader_manifest, + legacy_period_matches, + validate_legacy_market_loader_manifest, +) __all__ = [ @@ -168,8 +174,13 @@ def _parquet_file( metadata: Mapping[str, Any], policy: ExecutionPolicy, ) -> pq.ParquetFile: - if metadata.get("object_kind") != "PARQUET": - raise MarketDataValidationError("object_kind must be PARQUET") + expected_kind = ( + "MARKET_BARS" + if policy.market_data_schema_version == LEGACY_MARKET_SCHEMA_ID + else "PARQUET" + ) + if metadata.get("object_kind") != expected_kind: + raise MarketDataValidationError(f"object_kind must be {expected_kind}") if metadata.get("schema_version") != policy.market_data_schema_version: raise MarketDataValidationError("manifest schema_version does not match policy") path = self._object_path(metadata.get("object_key")) @@ -190,9 +201,13 @@ def _validate_manifest( manifest: Mapping[str, Any], policy: ExecutionPolicy, ) -> None: + legacy = is_legacy_market_loader_manifest(manifest) try: - validate_dataset_manifest(manifest) - except ContractValidationError as exc: + if legacy: + validate_legacy_market_loader_manifest(manifest) + else: + validate_dataset_manifest(manifest) + except (ContractValidationError, ValueError) as exc: raise MarketDataValidationError(str(exc)) from exc if self._manifest_validator is not None: try: @@ -210,16 +225,25 @@ def _validate_manifest( ) if manifest.get("schema_id") != policy.market_data_schema_version: raise MarketDataValidationError("manifest schema_id does not match policy") - if ( - _utc_timestamp(manifest.get("period_start"), "manifest.period_start") - != policy.period_start - ): - raise MarketDataValidationError("manifest period_start does not match policy") - if ( - _utc_timestamp(manifest.get("period_end"), "manifest.period_end") - != policy.period_end - ): - raise MarketDataValidationError("manifest period_end does not match policy") + if legacy: + if not legacy_period_matches( + manifest, + policy.period_start, + policy.period_end, + policy.timezone, + ): + raise MarketDataValidationError("legacy manifest period does not match policy") + else: + if ( + _utc_timestamp(manifest.get("period_start"), "manifest.period_start") + != policy.period_start + ): + raise MarketDataValidationError("manifest period_start does not match policy") + if ( + _utc_timestamp(manifest.get("period_end"), "manifest.period_end") + != policy.period_end + ): + raise MarketDataValidationError("manifest period_end does not match policy") def iter_batches( self, diff --git a/src/backtest_engine/production.py b/src/backtest_engine/production.py index b3745db..6ca5758 100644 --- a/src/backtest_engine/production.py +++ b/src/backtest_engine/production.py @@ -215,9 +215,11 @@ def by_id(self, manifest_id: uuid.UUID) -> Mapping[str, Any] | None: SELECT manifest.id, manifest.revision_number, manifest.status, manifest.dataset_hash, manifest.schema_version, manifest.period_start, manifest.period_end, manifest.available_at, - feed.resolution AS feed_resolution + provider.code AS provider_code, feed.code AS feed_code, + manifest.data_layer, feed.resolution AS feed_resolution FROM market_data.dataset_manifests manifest JOIN market_data.feeds feed ON feed.id = manifest.feed_id + JOIN market_data.providers provider ON provider.id = feed.provider_id WHERE manifest.id = :manifest_id """ ) @@ -226,7 +228,9 @@ def by_id(self, manifest_id: uuid.UUID) -> Mapping[str, Any] | None: SELECT o.id AS storage_object_id, o.object_key, o.content_hash, d.object_kind, d.partition_granularity, d.partition_start, d.partition_end, d.period_start, d.period_end, d.shard_key, - d.part_number, d.row_count, o.schema_version, o.status + d.part_number, d.row_count, o.schema_version, o.status, + o.storage_provider, o.bucket_name, o.provider_version_id, + o.file_format, o.media_type FROM market_data.dataset_objects d JOIN storage.objects o ON o.id = d.object_id WHERE d.dataset_manifest_id = :manifest_id @@ -246,6 +250,15 @@ def by_id(self, manifest_id: uuid.UUID) -> Mapping[str, Any] | None: raise _unavailable_manifest( f"dataset manifest {manifest_id} references a non-AVAILABLE object" ) + if any( + item["storage_provider"] != "S3" + or not item["bucket_name"] + or not item["provider_version_id"] + for item in object_rows + ): + raise _unavailable_manifest( + f"dataset manifest {manifest_id} lacks immutable S3 version evidence" + ) objects = [ { "storage_object_id": str(item["storage_object_id"]), @@ -261,6 +274,11 @@ def by_id(self, manifest_id: uuid.UUID) -> Mapping[str, Any] | None: "part_number": int(item["part_number"]), "row_count": int(item["row_count"]), "schema_version": str(item["schema_version"]), + "storage_provider": str(item["storage_provider"]), + "bucket_name": str(item["bucket_name"]), + "provider_version_id": str(item["provider_version_id"]), + "file_format": str(item["file_format"]), + "media_type": str(item["media_type"]), } for item in object_rows ] @@ -285,6 +303,9 @@ def by_id(self, manifest_id: uuid.UUID) -> Mapping[str, Any] | None: "status": str(row["status"]), "dataset_hash": str(row["dataset_hash"]), "schema_id": str(row["schema_version"]), + "provider_code": str(row["provider_code"]), + "feed_code": str(row["feed_code"]), + "data_layer": str(row["data_layer"]), "resolution": resolution, "period_start": _utc_text(row["period_start"]), "period_end": _utc_text(row["period_end"]), @@ -710,6 +731,15 @@ def materialize(self, manifest: Mapping[str, Any]) -> None: for metadata in manifest.get("objects", []): key = str(metadata.get("object_key", "")) expected = str(metadata.get("content_hash", "")) + storage_provider = str(metadata.get("storage_provider", "")) + bucket = str(metadata.get("bucket_name", "")) + version_id = str(metadata.get("provider_version_id", "")) + if storage_provider != "S3": + raise ConfigurationError(f"market-data object is not stored in S3: {key}") + if bucket != self._bucket: + raise ConfigurationError(f"market-data object bucket does not match configuration: {key}") + if not version_id: + raise ConfigurationError(f"market-data object has no immutable provider version: {key}") target = self._target(key) if target.is_file(): with target.open("rb") as existing: @@ -719,7 +749,7 @@ def materialize(self, manifest: Mapping[str, Any]) -> None: target.unlink() target.parent.mkdir(parents=True, exist_ok=True) temporary = target.with_name(f".{target.name}.{uuid.uuid4().hex}.part") - response = self._client.get_object(Bucket=self._bucket, Key=key) + response = self._client.get_object(Bucket=bucket, Key=key, VersionId=version_id) body = response["Body"] try: actual = self._copy_and_hash(body, temporary) diff --git a/src/backtest_engine/wiring.py b/src/backtest_engine/wiring.py index 4737ba9..735d9fa 100644 --- a/src/backtest_engine/wiring.py +++ b/src/backtest_engine/wiring.py @@ -114,6 +114,11 @@ FeatureOutputBindingError, resolve_feature_materialization_pins, ) +from .legacy_market_data import ( + is_legacy_market_loader_manifest, + legacy_period_matches, + validate_legacy_market_loader_manifest, +) from .lifecycle import BacktestLifecycleService, PersistenceRunGateway, SqsBacktestJobQueue from .money import PRECISION_RULES_VERSION, QUANTITY_QUANTUM, apply_rate, quantize_money, quantize_quantity from .monthly_judgment import ( @@ -1392,7 +1397,22 @@ def require_compatible_execution_window( problems: list[str] = [] manifest_start = _parse_instant(manifest.get("period_start")) manifest_end = _parse_instant(manifest.get("period_end")) - if manifest_start != policy.period_start or manifest_end != policy.period_end: + legacy = is_legacy_market_loader_manifest(manifest) + if legacy: + try: + validate_legacy_market_loader_manifest(manifest) + period_matches = legacy_period_matches( + manifest, + policy.period_start, + policy.period_end, + policy.timezone, + ) + except ValueError as exc: + problems.append(f"legacy dataset manifest is invalid: {exc}") + period_matches = False + else: + period_matches = manifest_start == policy.period_start and manifest_end == policy.period_end + if not period_matches: problems.append( "dataset manifest period " f"{manifest_start.isoformat()}..{manifest_end.isoformat()} does not match " diff --git a/tests/fixtures/int03-development-market-manifest.json b/tests/fixtures/int03-development-market-manifest.json new file mode 100644 index 0000000..9dca3da --- /dev/null +++ b/tests/fixtures/int03-development-market-manifest.json @@ -0,0 +1,79 @@ +{ + "contract_id": "com06.dataset-manifest", + "schema_version": 1, + "manifest_id": "7f7113c9-3b02-4098-97ec-0baa07e2b3b0", + "dataset_id": "7f7113c9-3b02-4098-97ec-0baa07e2b3b0", + "revision": 1, + "status": "AVAILABLE", + "dataset_hash": "08a848a5f9aa1aac80e215c2d86bcf6d5f96c354400c7c16394dae9ffa9939af", + "schema_id": "market-bars/1", + "provider_code": "ALPACA", + "feed_code": "ALPACA_SIP_ALL_30M", + "data_layer": "ADJUSTED", + "resolution": "30m", + "period_start": "2024-01-01T00:00:00Z", + "period_end": "2024-02-01T00:00:00Z", + "available_at": "2026-07-30T06:00:29Z", + "objects": [ + { + "storage_object_id": "de38b54e-d0fa-448d-bc56-fd11b21785af", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=00-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "256fc05e908c010e863d5ade3bff5776698a37d376605b2fbc2fc197203400d4", + "object_kind": "MARKET_BARS", + "partition_granularity": "YEAR", + "partition_start": "2024-01-01", + "partition_end": "2024-02-01", + "period_start": "2024-01-01T00:00:00Z", + "period_end": "2024-02-01T00:00:00Z", + "shard_key": "s00-of-8", + "part_number": 1, + "row_count": 0, + "schema_version": "market-bars/1", + "storage_provider": "S3", + "bucket_name": "idea2strategy-dev-418553863687-market-data", + "provider_version_id": "0jZPVZwe6fmVLMx9jT7sePpYAPETRnYs" + }, + { + "storage_object_id": "509f6d2f-ce0e-45d8-b611-3a0669ae62e5", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=01-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "3c82aada0b061ae7c10ff7d46574409c2b5bdc1f48c7531ff467547a9c6d6501", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s01-of-8", "part_number": 1, "row_count": 0, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "hz6laNWy7BrFenHJ6GNOM9t089BOBeLc" + }, + { + "storage_object_id": "b5f20f8e-9ced-427c-92b6-f40aabd91530", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=02-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "642897a31f11776bc6240db9aa6543bf046ee94eb14b40588df57e03ced6e215", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s02-of-8", "part_number": 1, "row_count": 0, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "TSekqXJZcU4vhtQPIR2_qa64Umjlp1mc" + }, + { + "storage_object_id": "d8aa63e2-e278-4b89-804f-0b66182f2210", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=03-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "379a8e4b065e116fcc5fa5ab0108ac6365b5fe60f5c38c6b4055f38c89478f98", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s03-of-8", "part_number": 1, "row_count": 0, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "p8gqQQYFHac9AAKpU9f0FIrPqgjHo5a_" + }, + { + "storage_object_id": "37d91ad1-36dc-45db-b611-3a0669ae62e5", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=04-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "216e4e0540acc13abb61a6e9f5e98703f3bd59fb7778b3513ad8c20841b11ea1", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s04-of-8", "part_number": 1, "row_count": 546, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "F9IO8ByKzDuFrF_EABssaa8mns0pDvPB" + }, + { + "storage_object_id": "cbd6445f-db3b-43f2-b8ea-9ed0f4351056", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=05-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "dc491bfaede62887742e370e1726d216fbf0e91d9826e063fc66384be1491720", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s05-of-8", "part_number": 1, "row_count": 0, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "JW0mnJvrNQzVUP9I9h1dU2fPKl.0AWGN" + }, + { + "storage_object_id": "23db3489-01ac-488b-8ed6-5a491d98b85a", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=06-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "c3f178e265101aa1c270f8361d32ffebb0ac98285ea7940da5e608a919d2b209", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s06-of-8", "part_number": 1, "row_count": 0, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "aTkdOrYbNNolS0HPg6YLMxSXRYuLo2nQ" + }, + { + "storage_object_id": "1fbc2ce0-33ae-4de6-a5f0-aff0ec10a99f", + "object_key": "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/resolution=30m/revision=00000001/year=2024/shard=07-of-08/manifest_id=7f7113c9-3b02-4098-97ec-0baa07e2b3b0/part-00001.parquet", + "content_hash": "2fbbdb9c5fceba3adf2705a2ed77ca1c451e3288c6e5df97866bf8dec9453866", + "object_kind": "MARKET_BARS", "partition_granularity": "YEAR", "partition_start": "2024-01-01", "partition_end": "2024-02-01", "period_start": "2024-01-01T00:00:00Z", "period_end": "2024-02-01T00:00:00Z", "shard_key": "s07-of-8", "part_number": 1, "row_count": 0, "schema_version": "market-bars/1", "storage_provider": "S3", "bucket_name": "idea2strategy-dev-418553863687-market-data", "provider_version_id": "5VxSqchVdkZQrha09_FdtrtdCP54mk9p" + } + ] +} diff --git a/tests/test_execution_policy.py b/tests/test_execution_policy.py index 0e3e3e4..333c10e 100644 --- a/tests/test_execution_policy.py +++ b/tests/test_execution_policy.py @@ -202,7 +202,7 @@ def test_catalog_rejects_unpublished_quarter_and_naive_release_time() -> None: def test_policy_rejects_an_intraday_period_presented_as_a_quarter() -> None: # The pre-rebuild D17 fixture declared 2024-Q1 with a 60-minute period. - with pytest.raises(ValueError, match="ET calendar quarter"): + with pytest.raises(ValueError, match="local calendar-day boundaries"): _policy( "2024-Q1", "official-backtest-policy-v1", @@ -211,14 +211,16 @@ def test_policy_rejects_an_intraday_period_presented_as_a_quarter() -> None: ) -def test_policy_rejects_a_period_that_ends_mid_quarter() -> None: - with pytest.raises(ValueError, match="ET calendar quarter"): - _policy( - "2024-Q1", - "official-backtest-policy-v1", - et_quarter_start(2024, 1), - datetime(2024, 2, 1, 5, 0, tzinfo=timezone.utc), - ) +def test_policy_accepts_a_complete_et_calendar_month() -> None: + policy = _policy( + "2024-Q1", + "official-backtest-policy-v1", + et_quarter_start(2024, 1), + datetime(2024, 2, 1, 5, 0, tzinfo=timezone.utc), + ) + + with pytest.raises(ValueError, match="not aligned to ET quarter"): + _ = policy.quarter_count def test_policy_rejects_a_reversed_period() -> None: diff --git a/tests/test_legacy_market_data.py b/tests/test_legacy_market_data.py new file mode 100644 index 0000000..191f224 --- /dev/null +++ b/tests/test_legacy_market_data.py @@ -0,0 +1,151 @@ +from __future__ import annotations + +import hashlib +import json +from copy import deepcopy +from dataclasses import replace +from datetime import datetime, timezone +from pathlib import Path + +import pyarrow.parquet as pq +import pytest + +from backtest_engine.calendar import XNYS_CALENDAR +from backtest_engine.execution_policy import D17_EXECUTION_POLICY_FIXTURE +from backtest_engine.legacy_market_data import ( + legacy_dataset_hash, + validate_legacy_market_loader_manifest, +) +from backtest_engine.market_data import MarketDataValidationError, ParquetMarketDataReader +from backtest_engine.wiring import require_compatible_execution_window +from d_market_data_testkit import write_small_market_bars + + +FIXTURE = Path(__file__).parent / "fixtures" / "int03-development-market-manifest.json" +PERIOD_START = "2024-01-01T00:00:00Z" +PERIOD_END = "2024-02-01T00:00:00Z" + + +def _fixture() -> dict[str, object]: + return json.loads(FIXTURE.read_text(encoding="utf-8")) + + +def test_exact_development_manifest_recomputes_the_legacy_loader_hash() -> None: + manifest = _fixture() + + validate_legacy_market_loader_manifest(manifest) + + assert legacy_dataset_hash(manifest) == ("08a848a5f9aa1aac80e215c2d86bcf6d5f96c354400c7c16394dae9ffa9939af") + assert sum(item["row_count"] for item in manifest["objects"]) == 546 # type: ignore[index,union-attr] + + +def test_exact_development_manifest_binds_to_the_one_month_et_policy() -> None: + policy = replace( + D17_EXECUTION_POLICY_FIXTURE, + version="development-official-backtest-2026-q3-v2", + period_start=datetime(2024, 1, 1, 5, tzinfo=timezone.utc), + period_end=datetime(2024, 2, 1, 5, tzinfo=timezone.utc), + market_data_schema_version="market-bars/1", + ) + + require_compatible_execution_window(policy, _fixture(), XNYS_CALENDAR) + + +@pytest.mark.parametrize( + "mutate", + [ + lambda document: document["objects"][4].update(row_count=545), + lambda document: document["objects"][4].update(provider_version_id=""), + lambda document: document["objects"][4].update(object_kind="PARQUET"), + lambda document: document["objects"][4].update(shard_key="s03-of-8"), + lambda document: document.update(data_layer="RAW"), + ], +) +def test_legacy_adapter_rejects_tampered_catalog_evidence(mutate: object) -> None: + manifest = deepcopy(_fixture()) + mutate(manifest) # type: ignore[operator] + + with pytest.raises(ValueError): + validate_legacy_market_loader_manifest(manifest) + + +def _one_shard_manifest(path: Path) -> dict[str, object]: + key = ( + "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/" + "resolution=30m/revision=00000001/year=2024/shard=00-of-01/" + "manifest_id=11111111-1111-4111-8111-111111111111/part-00001.parquet" + ) + metadata = { + "storage_object_id": "22222222-2222-4222-8222-222222222222", + "object_key": key, + "content_hash": hashlib.sha256(path.read_bytes()).hexdigest(), + "object_kind": "MARKET_BARS", + "partition_granularity": "YEAR", + "partition_start": "2024-01-01", + "partition_end": "2024-02-01", + "period_start": PERIOD_START, + "period_end": PERIOD_END, + "shard_key": "s00-of-1", + "part_number": 1, + "row_count": 2, + "schema_version": "market-bars/1", + } + manifest: dict[str, object] = { + "contract_id": "com06.dataset-manifest", + "schema_version": 1, + "manifest_id": "11111111-1111-4111-8111-111111111111", + "dataset_id": "11111111-1111-4111-8111-111111111111", + "revision": 1, + "status": "AVAILABLE", + "dataset_hash": "", + "schema_id": "market-bars/1", + "provider_code": "ALPACA", + "feed_code": "ALPACA_SIP_ALL_30M", + "data_layer": "ADJUSTED", + "resolution": "30m", + "period_start": PERIOD_START, + "period_end": PERIOD_END, + "available_at": "2026-07-30T06:00:29Z", + "objects": [metadata], + } + manifest["dataset_hash"] = legacy_dataset_hash(manifest) + return manifest + + +def test_reader_consumes_legacy_parquet_without_weakening_the_canonical_path( + tmp_path: Path, +) -> None: + key_root = ( + tmp_path + / "historical/provider=alpaca/feed=sip/adjustment=all/session=regular" + / "resolution=30m/revision=00000001/year=2024/shard=00-of-01" + / "manifest_id=11111111-1111-4111-8111-111111111111" + ) + key_root.mkdir(parents=True) + fixture = write_small_market_bars(key_root / "part-00001.parquet") + table = pq.read_table(fixture.path).replace_schema_metadata( + {b"schema_version": b"market-bars/1", b"processing_version": b"market-loader/1.0.0"} + ) + pq.write_table(table, fixture.path, compression="zstd", version="2.6") + manifest = _one_shard_manifest(fixture.path) + policy = replace( + D17_EXECUTION_POLICY_FIXTURE, + version="development-official-backtest-2026-q3-v2", + period_start=datetime(2024, 1, 1, 5, tzinfo=timezone.utc), + period_end=datetime(2024, 2, 1, 5, tzinfo=timezone.utc), + market_data_schema_version="market-bars/1", + ) + + result = ParquetMarketDataReader(tmp_path).read(manifest, policy) + + assert result.num_rows == 2 + + +def test_canonical_manifest_cannot_claim_the_legacy_object_kind(tmp_path: Path) -> None: + fixture = write_small_market_bars(tmp_path / "market-bars.parquet") + manifest = _one_shard_manifest(fixture.path) + manifest["schema_id"] = "market-bars-v2" + manifest["objects"][0]["partition_granularity"] = "DAY" # type: ignore[index] + + with pytest.raises(MarketDataValidationError, match="object_kind"): + ParquetMarketDataReader(tmp_path).read(manifest, D17_EXECUTION_POLICY_FIXTURE) diff --git a/tests/test_production.py b/tests/test_production.py index 8a4cc9a..70db3f9 100644 --- a/tests/test_production.py +++ b/tests/test_production.py @@ -260,6 +260,8 @@ def test_compiled_plan_source_returns_the_immutable_launch_contract_document() - def _dataset_manifest_source( manifest_id: UUID, object_keys: list[str], + *, + object_overrides: dict[str, object] | None = None, ) -> PostgresDatasetManifestSource: manifest = { "id": manifest_id, @@ -270,6 +272,9 @@ def _dataset_manifest_source( "period_start": datetime(2024, 1, 1, tzinfo=UTC), "period_end": datetime(2024, 2, 1, tzinfo=UTC), "available_at": datetime(2026, 8, 9, tzinfo=UTC), + "provider_code": "ALPACA", + "feed_code": "ALPACA_SIP_ALL_30M", + "data_layer": "ADJUSTED", "feed_resolution": "30m", } objects = [ @@ -287,10 +292,17 @@ def _dataset_manifest_source( "part_number": 1, "row_count": 10, "schema_version": "market-bars/1", + "storage_provider": "S3", + "bucket_name": "market", + "provider_version_id": f"version-{index}", + "file_format": "PARQUET", + "media_type": "application/vnd.apache.parquet", "status": "AVAILABLE", } for index, object_key in enumerate(object_keys, start=1) ] + for item in objects: + item.update(object_overrides or {}) return PostgresDatasetManifestSource(_Engine(manifest, objects)) # type: ignore[arg-type] @@ -311,6 +323,26 @@ def test_dataset_manifest_source_accepts_the_deployed_legacy_loader_binding() -> assert resolved["manifest_id"] == str(manifest_id) assert resolved["dataset_id"] == str(manifest_id) assert [item["object_key"] for item in resolved["objects"]] == object_keys + assert [item["provider_version_id"] for item in resolved["objects"]] == [ + f"version-{index}" for index in range(1, 9) + ] + + +def test_dataset_manifest_source_rejects_an_object_without_an_immutable_s3_version() -> None: + manifest_id = UUID("7f7113c9-3b02-4098-97ec-0baa07e2b3b0") + key = ( + "historical/provider=alpaca/feed=sip/adjustment=all/session=regular/" + "resolution=30m/revision=00000001/year=2024/shard=00-of-01/" + f"manifest_id={manifest_id}/part-00001.parquet" + ) + source = _dataset_manifest_source( + manifest_id, + [key], + object_overrides={"provider_version_id": None}, + ) + + with pytest.raises(JobNotSatisfiable, match="immutable S3 version evidence"): + source.by_id(manifest_id) def test_dataset_manifest_source_preserves_the_canonical_logical_dataset_binding() -> None: @@ -614,7 +646,11 @@ def close(self) -> None: class _S3: def get_object(self, **kwargs: str) -> dict[str, object]: - assert kwargs == {"Bucket": "market", "Key": "year/part.parquet"} + assert kwargs == { + "Bucket": "market", + "Key": "year/part.parquet", + "VersionId": "immutable-version-1", + } return {"Body": _Body()} reader = S3ParquetMarketDataReader(bucket="market", cache_root=tmp_path, client=_S3()) @@ -623,6 +659,9 @@ def get_object(self, **kwargs: str) -> dict[str, object]: { "object_key": "year/part.parquet", "content_hash": hashlib.sha256(body).hexdigest(), + "storage_provider": "S3", + "bucket_name": "market", + "provider_version_id": "immutable-version-1", } ] } @@ -647,5 +686,43 @@ def close(self) -> None: reader = S3ParquetMarketDataReader(bucket="market", cache_root=tmp_path, client=client) with pytest.raises(ConfigurationError, match="checksum"): - reader.materialize({"objects": [{"object_key": "part.parquet", "content_hash": "a" * 64}]}) + reader.materialize( + { + "objects": [ + { + "object_key": "part.parquet", + "content_hash": "a" * 64, + "storage_provider": "S3", + "bucket_name": "market", + "provider_version_id": "immutable-version-1", + } + ] + } + ) assert not (tmp_path / "part.parquet").exists() + + +@pytest.mark.parametrize( + ("field", "value"), + [ + ("storage_provider", "LOCAL"), + ("bucket_name", "other-market"), + ("provider_version_id", ""), + ], +) +def test_s3_reader_rejects_mutable_or_substituted_object_evidence( + tmp_path: Path, field: str, value: str +) -> None: + metadata = { + "object_key": "part.parquet", + "content_hash": "a" * 64, + "storage_provider": "S3", + "bucket_name": "market", + "provider_version_id": "immutable-version-1", + } + metadata[field] = value + client = SimpleNamespace(get_object=lambda **_kwargs: pytest.fail("must fail before S3")) + reader = S3ParquetMarketDataReader(bucket="market", cache_root=tmp_path, client=client) + + with pytest.raises(ConfigurationError): + reader.materialize({"objects": [metadata]})