2524 lines
117 KiB
Python
2524 lines
117 KiB
Python
"""Immutable, versioned factor-definition and factor-set computation contracts.
|
|
|
|
The module is deliberately storage and provider neutral. It validates complete
|
|
DatasetSnapshot and Data Foundation v1 envelopes at the consumer boundary, but
|
|
does not fetch records, execute factor formulas, persist artifacts, or grant
|
|
decision, production, or live eligibility.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import re
|
|
from collections import defaultdict
|
|
from collections.abc import Mapping, Sequence
|
|
from dataclasses import dataclass, field, fields
|
|
from datetime import datetime
|
|
from enum import StrEnum
|
|
from itertools import pairwise
|
|
from types import MappingProxyType
|
|
from typing import Any, Final, Never, Self, cast
|
|
|
|
from quant_engine.alpha_factors import ALPHA158_REGISTRY
|
|
|
|
MAX_SAFE_INTEGER: Final = 9_007_199_254_740_991
|
|
FACTOR_CONTRACT_VERSION: Final = "1.0.0"
|
|
|
|
_CANONICAL_KEY = re.compile(r"^[a-z][a-z0-9_]*$")
|
|
_FIELD_NAME = re.compile(
|
|
r"^(?!.*(?:latest|provider|table|sql|locator|source|tushare|qtdb|edb))"
|
|
r"[a-z][a-z0-9]*(?:_[a-z0-9]+)*$"
|
|
)
|
|
_LOGICAL_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]{0,127}$")
|
|
_SEMVER = re.compile(r"^(?:0|[1-9][0-9]*)\.(?:0|[1-9][0-9]*)\.(?:0|[1-9][0-9]*)$")
|
|
_SHA256 = re.compile(r"^sha256:[0-9a-f]{64}$")
|
|
_BARE_SHA256 = re.compile(r"^[0-9a-f]{64}$")
|
|
_GIT_SHA = re.compile(r"^[0-9a-f]{40}$")
|
|
_UTC_INSTANT = re.compile(
|
|
r"^[0-9]{4}-(?:0[1-9]|1[0-2])-(?:0[1-9]|[12][0-9]|3[01])"
|
|
r"T(?:[01][0-9]|2[0-3]):[0-5][0-9]:[0-5][0-9](?:\.[0-9]{1,6})?Z$"
|
|
)
|
|
_CALENDAR_DATE = re.compile(
|
|
r"^[0-9]{4}-(?:0[1-9]|1[0-2])-(?:0[1-9]|[12][0-9]|3[01])$"
|
|
)
|
|
_CANONICAL_DECIMAL = re.compile(r"^-?(?:0|[1-9][0-9]*)(?:\.[0-9]*[1-9])?$")
|
|
_DATASET_ID = re.compile(r"^rhdataset:(?:market|macroeconomic):[0-9a-f]{32}$")
|
|
_TRANSFORMATION_ID = re.compile(r"^rhtransform:[0-9a-f]{32}$")
|
|
_SNAPSHOT_ID = re.compile(r"^rhdsv1:sha256:[0-9a-f]{64}$")
|
|
_FOUNDATION_ID = re.compile(r"^rhdfv1:sha256:[0-9a-f]{64}$")
|
|
_INSTRUMENT_ID = re.compile(r"^rhinstrument:[0-9a-f]{32}$")
|
|
_ROUTE_ID = re.compile(r"^rhroutev1:sha256:[0-9a-f]{64}$")
|
|
_CALENDAR_ID = re.compile(r"^rhcalendar:[0-9a-f]{32}$")
|
|
_CALENDAR_REVISION_ID = re.compile(r"^rhcalv1:sha256:[0-9a-f]{64}$")
|
|
_ACTION_ID = re.compile(r"^rhaction:[0-9a-f]{32}$")
|
|
_ACTION_REVISION_ID = re.compile(r"^rhcav1:sha256:[0-9a-f]{64}$")
|
|
_VIEW_ID = re.compile(r"^rhview:[0-9a-f]{32}$")
|
|
_VIEW_REF_ID = re.compile(r"^rhviewrefv1:sha256:[0-9a-f]{64}$")
|
|
_DEFINITION_ID = re.compile(r"^rhfactorv1:sha256:[0-9a-f]{64}$")
|
|
_FACTOR_SET_ID = re.compile(r"^rhfactorsetv1:sha256:[0-9a-f]{64}$")
|
|
_OUTPUT_ARTIFACT_ID = re.compile(r"^rhfactoroutputv1:sha256:[0-9a-f]{64}$")
|
|
_LEGACY_BINDING_ID = re.compile(r"^rhlegacyfactorv1:sha256:[0-9a-f]{64}$")
|
|
_FORBIDDEN_TOKENS: Final = frozenset(
|
|
{"latest", "provider", "table", "sql", "locator", "tushare", "wind", "bloomberg"}
|
|
)
|
|
|
|
|
|
class ContractErrorCode(StrEnum):
|
|
"""Stable machine-readable rejection categories for the public boundary."""
|
|
|
|
TYPE_ERROR = "type_error"
|
|
MISSING_FIELD = "missing_field"
|
|
UNKNOWN_FIELD = "unknown_field"
|
|
INVALID_FORMAT = "invalid_format"
|
|
INVALID_VALUE = "invalid_value"
|
|
IDENTITY_MISMATCH = "identity_mismatch"
|
|
QUALIFICATION_REJECTED = "qualification_rejected"
|
|
TIME_ORDER_VIOLATION = "time_order_violation"
|
|
INPUT_CLOSURE_VIOLATION = "input_closure_violation"
|
|
LINEAGE_VIOLATION = "lineage_violation"
|
|
ARTIFACT_MISMATCH = "artifact_mismatch"
|
|
READINESS_ESCALATION = "readiness_escalation"
|
|
LEGACY_BINDING_MISMATCH = "legacy_binding_mismatch"
|
|
|
|
|
|
class FactorContractError(ValueError):
|
|
"""Typed deterministic contract rejection with an invariant/JSON path."""
|
|
|
|
def __init__(self, code: ContractErrorCode, path: str, detail: str) -> None:
|
|
self.code = code
|
|
self.path = path
|
|
self.detail = detail
|
|
super().__init__(f"{code.value} at {path}: {detail}")
|
|
|
|
|
|
def _fail(code: ContractErrorCode, path: str, detail: str) -> Never:
|
|
raise FactorContractError(code, path, detail)
|
|
|
|
|
|
def _assert_canonical_profile(value: Any, path: str = "$") -> None:
|
|
if type(value) is dict:
|
|
for key, child in value.items():
|
|
if type(key) is not str or _CANONICAL_KEY.fullmatch(key) is None:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, path, "non-canonical object key")
|
|
_assert_canonical_profile(child, f"{path}.{key}")
|
|
return
|
|
if type(value) is list:
|
|
for index, child in enumerate(value):
|
|
_assert_canonical_profile(child, f"{path}[{index}]")
|
|
return
|
|
if value is None or type(value) is bool:
|
|
return
|
|
if type(value) is int:
|
|
if not -MAX_SAFE_INTEGER <= value <= MAX_SAFE_INTEGER:
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, "integer exceeds safe range")
|
|
return
|
|
if type(value) is str:
|
|
try:
|
|
value.encode("utf-8")
|
|
except UnicodeEncodeError as exc:
|
|
raise FactorContractError(
|
|
ContractErrorCode.INVALID_FORMAT,
|
|
path,
|
|
"string contains an unpaired surrogate",
|
|
) from exc
|
|
return
|
|
_fail(
|
|
ContractErrorCode.TYPE_ERROR,
|
|
path,
|
|
"value is outside the restricted canonical JSON domain",
|
|
)
|
|
|
|
|
|
def canonical_json_bytes(value: Any) -> bytes:
|
|
"""Return bytes for the restricted ResearchHub canonical JSON profile."""
|
|
|
|
_assert_canonical_profile(value)
|
|
return json.dumps(
|
|
value,
|
|
allow_nan=False,
|
|
ensure_ascii=False,
|
|
separators=(",", ":"),
|
|
sort_keys=True,
|
|
).encode("utf-8")
|
|
|
|
|
|
def canonical_json(value: Any) -> str:
|
|
return canonical_json_bytes(value).decode("utf-8")
|
|
|
|
|
|
def _duplicate_key_pairs(pairs: list[tuple[str, Any]]) -> dict[str, Any]:
|
|
result: dict[str, Any] = {}
|
|
for key, value in pairs:
|
|
if key in result:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$", f"duplicate JSON key: {key}")
|
|
result[key] = value
|
|
return result
|
|
|
|
|
|
def _parse_json_object(value: str | bytes, path: str) -> dict[str, Any]:
|
|
try:
|
|
source_bytes = value if isinstance(value, bytes) else value.encode("utf-8")
|
|
except UnicodeEncodeError as exc:
|
|
raise FactorContractError(
|
|
ContractErrorCode.INVALID_FORMAT,
|
|
path,
|
|
"invalid UTF-8 JSON",
|
|
) from exc
|
|
try:
|
|
loaded = json.loads(value, object_pairs_hook=_duplicate_key_pairs)
|
|
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
|
|
raise FactorContractError(
|
|
ContractErrorCode.INVALID_FORMAT,
|
|
path,
|
|
"invalid UTF-8 JSON",
|
|
) from exc
|
|
if type(loaded) is not dict:
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "top-level JSON value must be an object")
|
|
_assert_canonical_profile(loaded, path)
|
|
if source_bytes != canonical_json_bytes(loaded):
|
|
_fail(ContractErrorCode.INVALID_FORMAT, path, "JSON is not canonical")
|
|
return cast(dict[str, Any], loaded)
|
|
|
|
|
|
def _freeze_json(value: Any, path: str = "$") -> Any:
|
|
if type(value) is dict:
|
|
copied: dict[str, Any] = {}
|
|
for key, child in value.items():
|
|
if type(key) is not str or _CANONICAL_KEY.fullmatch(key) is None:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, path, "non-canonical object key")
|
|
copied[key] = _freeze_json(child, f"{path}.{key}")
|
|
return MappingProxyType(copied)
|
|
if type(value) in {list, tuple}:
|
|
return tuple(_freeze_json(child, f"{path}[{index}]") for index, child in enumerate(value))
|
|
if value is None or type(value) in {str, bool, int}:
|
|
_assert_canonical_profile(value, path)
|
|
return value
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "unsupported mutable or custom value")
|
|
|
|
|
|
def _thaw_json(value: Any) -> Any:
|
|
if isinstance(value, Mapping):
|
|
return {key: _thaw_json(child) for key, child in value.items()}
|
|
if type(value) is tuple:
|
|
return [_thaw_json(child) for child in value]
|
|
return cast(dict[str, Any], value)
|
|
|
|
|
|
def _content_address(document: dict[str, Any], identity_field: str, prefix: str) -> str:
|
|
payload = {key: value for key, value in document.items() if key != identity_field}
|
|
return f"{prefix}{hashlib.sha256(canonical_json_bytes(payload)).hexdigest()}"
|
|
|
|
|
|
def _digest_bytes(value: bytes) -> str:
|
|
return f"sha256:{hashlib.sha256(value).hexdigest()}"
|
|
|
|
|
|
def _object(
|
|
value: Any,
|
|
path: str,
|
|
required: Sequence[str],
|
|
optional: Sequence[str] = (),
|
|
) -> dict[str, Any]:
|
|
if type(value) is not dict:
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "must be an object")
|
|
missing = [key for key in required if key not in value]
|
|
if missing:
|
|
_fail(ContractErrorCode.MISSING_FIELD, f"{path}.{missing[0]}", "field is required")
|
|
unexpected = sorted(set(value) - set(required) - set(optional))
|
|
if unexpected:
|
|
_fail(
|
|
ContractErrorCode.UNKNOWN_FIELD,
|
|
f"{path}.{unexpected[0]}",
|
|
"field is not allowed by the closed-world contract",
|
|
)
|
|
return cast(dict[str, Any], value)
|
|
|
|
|
|
def _array(value: Any, path: str, *, minimum: int = 0, unique: bool = False) -> list[Any]:
|
|
if type(value) is not list:
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "must be an array")
|
|
if len(value) < minimum:
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, f"requires at least {minimum} item(s)")
|
|
if unique:
|
|
encoded = [canonical_json_bytes(item) for item in value]
|
|
if len(encoded) != len(set(encoded)):
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, "items must be unique")
|
|
return value
|
|
|
|
|
|
def _string(value: Any, path: str, pattern: re.Pattern[str] | None = None) -> str:
|
|
if type(value) is not str:
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "must be a string")
|
|
if pattern is not None and pattern.fullmatch(value) is None:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, path, "does not match the required profile")
|
|
return value
|
|
|
|
|
|
def _enum(value: Any, path: str, allowed: set[str]) -> str:
|
|
normalized = _string(value, path)
|
|
if normalized not in allowed:
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, "unsupported enum value")
|
|
return normalized
|
|
|
|
|
|
def _safe_integer(value: Any, path: str, *, minimum: int | None = None) -> int:
|
|
if type(value) is not int:
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "must be an integer")
|
|
if not -MAX_SAFE_INTEGER <= value <= MAX_SAFE_INTEGER:
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, "integer exceeds safe range")
|
|
if minimum is not None and value < minimum:
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, f"must be at least {minimum}")
|
|
return value
|
|
|
|
|
|
def _parse_utc(value: Any, path: str) -> datetime:
|
|
text = _string(value, path, _UTC_INSTANT)
|
|
try:
|
|
return datetime.fromisoformat(f"{text[:-1]}+00:00")
|
|
except ValueError as exc:
|
|
raise FactorContractError(
|
|
ContractErrorCode.INVALID_FORMAT,
|
|
path,
|
|
"not a real calendar instant",
|
|
) from exc
|
|
|
|
|
|
def _normalize_new_instant(value: Any, path: str) -> str:
|
|
parsed = _parse_utc(value, path)
|
|
if parsed.microsecond:
|
|
fraction = f"{parsed.microsecond:06d}".rstrip("0")
|
|
return parsed.strftime("%Y-%m-%dT%H:%M:%S") + f".{fraction}Z"
|
|
return parsed.strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
|
|
|
|
def _parse_date(value: Any, path: str) -> None:
|
|
text = _string(value, path, _CALENDAR_DATE)
|
|
try:
|
|
datetime.strptime(text, "%Y-%m-%d")
|
|
except ValueError as exc:
|
|
raise FactorContractError(
|
|
ContractErrorCode.INVALID_FORMAT,
|
|
path,
|
|
"not a real calendar date",
|
|
) from exc
|
|
|
|
|
|
def _logical_id(value: Any, path: str) -> str:
|
|
text = _string(value, path, _LOGICAL_ID)
|
|
tokens = {token.lower() for token in re.findall(r"[A-Za-z0-9]+", text)}
|
|
if not tokens.isdisjoint(_FORBIDDEN_TOKENS) or "://" in text:
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, "mutable alias or physical locator is forbidden")
|
|
return text
|
|
|
|
|
|
def _semver(value: Any, path: str) -> str:
|
|
return _string(value, path, _SEMVER)
|
|
|
|
|
|
def _digest(value: Any, path: str) -> str:
|
|
return _string(value, path, _SHA256)
|
|
|
|
|
|
def _git_revision(value: Any, path: str) -> str:
|
|
return _string(value, path, _GIT_SHA)
|
|
|
|
|
|
def _no_physical_leakage(value: Any, path: str = "$") -> None:
|
|
if type(value) is dict:
|
|
for key, child in value.items():
|
|
tokens = set(re.findall(r"[a-z0-9]+", key.lower()))
|
|
if not tokens.isdisjoint(_FORBIDDEN_TOKENS):
|
|
_fail(ContractErrorCode.INVALID_VALUE, f"{path}.{key}", "physical token leakage")
|
|
_no_physical_leakage(child, f"{path}.{key}")
|
|
elif type(value) is list:
|
|
for index, child in enumerate(value):
|
|
_no_physical_leakage(child, f"{path}[{index}]")
|
|
elif type(value) is str:
|
|
tokens = set(re.findall(r"[a-z0-9]+", value.lower()))
|
|
if not tokens.isdisjoint(_FORBIDDEN_TOKENS):
|
|
_fail(ContractErrorCode.INVALID_VALUE, path, "physical token leakage")
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ProducerIdentity:
|
|
id: str
|
|
version: str
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "id", _logical_id(self.id, "$.producer.id"))
|
|
object.__setattr__(self, "version", _semver(self.version, "$.producer.version"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {"id": self.id, "version": self.version}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.producer") -> Self:
|
|
item = _object(value, path, ("id", "version"))
|
|
return cls(
|
|
id=_logical_id(item["id"], f"{path}.id"),
|
|
version=_semver(item["version"], f"{path}.version"),
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class TypedParameter:
|
|
type: str
|
|
value: Any
|
|
|
|
def __post_init__(self) -> None:
|
|
parameter_type = _enum(
|
|
self.type,
|
|
"$.parameters.*.type",
|
|
{"boolean", "decimal", "integer", "json", "null", "string"},
|
|
)
|
|
raw = self.value
|
|
if parameter_type == "boolean" and type(raw) is not bool:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.parameters.*.value", "must be boolean")
|
|
if parameter_type == "integer":
|
|
_safe_integer(raw, "$.parameters.*.value")
|
|
if parameter_type == "string" and type(raw) is not str:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.parameters.*.value", "must be string")
|
|
if parameter_type == "null" and raw is not None:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.parameters.*.value", "must be null")
|
|
if parameter_type == "decimal":
|
|
decimal = _string(raw, "$.parameters.*.value", _CANONICAL_DECIMAL)
|
|
if decimal in {"-0", "-0.0"}:
|
|
_fail(
|
|
ContractErrorCode.INVALID_FORMAT,
|
|
"$.parameters.*.value",
|
|
"negative zero is forbidden",
|
|
)
|
|
if parameter_type == "json":
|
|
_assert_canonical_profile(_thaw_json(_freeze_json(raw)), "$.parameters.*.value")
|
|
object.__setattr__(self, "type", parameter_type)
|
|
object.__setattr__(self, "value", _freeze_json(raw, "$.parameters.*.value"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {"type": self.type, "value": _thaw_json(self.value)}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str) -> Self:
|
|
item = _object(value, path, ("type", "value"))
|
|
return cls(type=_string(item["type"], f"{path}.type"), value=item["value"])
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class FactorInput:
|
|
input_name: str
|
|
schema_digest: str
|
|
required_columns: tuple[str, ...]
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "input_name", _string(self.input_name, "$.inputs[].input_name", _FIELD_NAME))
|
|
object.__setattr__(self, "schema_digest", _digest(self.schema_digest, "$.inputs[].schema_digest"))
|
|
if type(self.required_columns) not in {tuple, list}:
|
|
_fail(
|
|
ContractErrorCode.TYPE_ERROR,
|
|
"$.inputs[].required_columns",
|
|
"must be a sequence",
|
|
)
|
|
columns = tuple(
|
|
_string(column, f"$.inputs[].required_columns[{index}]", _FIELD_NAME)
|
|
for index, column in enumerate(self.required_columns)
|
|
)
|
|
if not columns or len(columns) != len(set(columns)):
|
|
_fail(
|
|
ContractErrorCode.INVALID_VALUE,
|
|
"$.inputs[].required_columns",
|
|
"must be non-empty and unique",
|
|
)
|
|
object.__setattr__(self, "required_columns", columns)
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"input_name": self.input_name,
|
|
"schema_digest": self.schema_digest,
|
|
"required_columns": list(self.required_columns),
|
|
}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str) -> Self:
|
|
item = _object(value, path, ("input_name", "schema_digest", "required_columns"))
|
|
columns = _array(item["required_columns"], f"{path}.required_columns", minimum=1, unique=True)
|
|
return cls(
|
|
input_name=_string(item["input_name"], f"{path}.input_name"),
|
|
schema_digest=_string(item["schema_digest"], f"{path}.schema_digest"),
|
|
required_columns=tuple(columns),
|
|
)
|
|
|
|
|
|
def factor_input_schema_digest(inputs: Sequence[FactorInput]) -> str:
|
|
normalized = _normalize_factor_inputs(inputs)
|
|
payload = [item.to_dict() for item in normalized]
|
|
return _digest_bytes(canonical_json_bytes(payload))
|
|
|
|
|
|
def _normalize_factor_inputs(inputs: Sequence[FactorInput]) -> tuple[FactorInput, ...]:
|
|
if type(inputs) not in {tuple, list}:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.inputs", "must be a sequence")
|
|
normalized: list[FactorInput] = []
|
|
for index, item in enumerate(inputs):
|
|
if not isinstance(item, FactorInput):
|
|
_fail(ContractErrorCode.TYPE_ERROR, f"$.inputs[{index}]", "must be FactorInput")
|
|
normalized.append(item)
|
|
if not normalized:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.inputs", "at least one input is required")
|
|
names = [item.input_name for item in normalized]
|
|
if len(names) != len(set(names)):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.inputs", "duplicate input_name")
|
|
return tuple(sorted(normalized, key=lambda item: item.input_name))
|
|
|
|
|
|
@dataclass(frozen=True, slots=True, init=False)
|
|
class FactorDefinition:
|
|
contract_name: str
|
|
schema_version: str
|
|
definition_id: str
|
|
factor_id: str
|
|
version: str
|
|
formula: str
|
|
parameters: Mapping[str, TypedParameter]
|
|
implementation_digest: str
|
|
input_schema_digest: str
|
|
inputs: tuple[FactorInput, ...]
|
|
valid_from: str
|
|
valid_until: str
|
|
warmup_sessions: int
|
|
lag_sessions: int
|
|
producer: ProducerIdentity
|
|
code_revision: str
|
|
|
|
@classmethod
|
|
def create(
|
|
cls,
|
|
*,
|
|
factor_id: str,
|
|
version: str,
|
|
formula: str,
|
|
parameters: Mapping[str, TypedParameter],
|
|
implementation_digest: str,
|
|
input_schema_digest: str,
|
|
inputs: Sequence[FactorInput],
|
|
valid_from: str,
|
|
valid_until: str,
|
|
warmup_sessions: int,
|
|
lag_sessions: int,
|
|
producer: ProducerIdentity,
|
|
code_revision: str,
|
|
) -> Self:
|
|
return cls._build(
|
|
factor_id=factor_id,
|
|
version=version,
|
|
formula=formula,
|
|
parameters=parameters,
|
|
implementation_digest=implementation_digest,
|
|
input_schema_digest=input_schema_digest,
|
|
inputs=inputs,
|
|
valid_from=valid_from,
|
|
valid_until=valid_until,
|
|
warmup_sessions=warmup_sessions,
|
|
lag_sessions=lag_sessions,
|
|
producer=producer,
|
|
code_revision=code_revision,
|
|
supplied_definition_id=None,
|
|
)
|
|
|
|
@classmethod
|
|
def _build(
|
|
cls,
|
|
*,
|
|
factor_id: Any,
|
|
version: Any,
|
|
formula: Any,
|
|
parameters: Any,
|
|
implementation_digest: Any,
|
|
input_schema_digest: Any,
|
|
inputs: Sequence[FactorInput],
|
|
valid_from: Any,
|
|
valid_until: Any,
|
|
warmup_sessions: Any,
|
|
lag_sessions: Any,
|
|
producer: Any,
|
|
code_revision: Any,
|
|
supplied_definition_id: Any,
|
|
) -> Self:
|
|
normalized_factor_id = _logical_id(factor_id, "$.factor_id")
|
|
normalized_version = _semver(version, "$.version")
|
|
normalized_formula = _string(formula, "$.formula")
|
|
if not 1 <= len(normalized_formula) <= 4096:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.formula", "length must be 1..4096")
|
|
if type(parameters) is not dict:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.parameters", "must be an object")
|
|
normalized_parameters: dict[str, TypedParameter] = {}
|
|
for name, value in parameters.items():
|
|
normalized_name = _string(name, "$.parameters", _FIELD_NAME)
|
|
if not isinstance(value, TypedParameter):
|
|
_fail(
|
|
ContractErrorCode.TYPE_ERROR,
|
|
f"$.parameters.{normalized_name}",
|
|
"must be TypedParameter",
|
|
)
|
|
normalized_parameters[normalized_name] = TypedParameter(value.type, _thaw_json(value.value))
|
|
normalized_inputs = _normalize_factor_inputs(inputs)
|
|
normalized_input_digest = _digest(input_schema_digest, "$.input_schema_digest")
|
|
expected_input_digest = factor_input_schema_digest(normalized_inputs)
|
|
if normalized_input_digest != expected_input_digest:
|
|
_fail(
|
|
ContractErrorCode.IDENTITY_MISMATCH,
|
|
"$.input_schema_digest",
|
|
"does not identify the canonical input declarations",
|
|
)
|
|
normalized_valid_from = _normalize_new_instant(valid_from, "$.valid_from")
|
|
normalized_valid_until = _normalize_new_instant(valid_until, "$.valid_until")
|
|
if _parse_utc(normalized_valid_from, "$.valid_from") >= _parse_utc(
|
|
normalized_valid_until, "$.valid_until"
|
|
):
|
|
_fail(
|
|
ContractErrorCode.TIME_ORDER_VIOLATION,
|
|
"$.validity",
|
|
"valid_from must be before valid_until",
|
|
)
|
|
normalized_warmup = _safe_integer(warmup_sessions, "$.warmup_sessions", minimum=0)
|
|
normalized_lag = _safe_integer(lag_sessions, "$.lag_sessions", minimum=0)
|
|
if not isinstance(producer, ProducerIdentity):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.producer", "must be ProducerIdentity")
|
|
payload = {
|
|
"contract_name": "researchhub.factor-definition",
|
|
"schema_version": FACTOR_CONTRACT_VERSION,
|
|
"factor_id": normalized_factor_id,
|
|
"version": normalized_version,
|
|
"formula": normalized_formula,
|
|
"parameters": {
|
|
name: normalized_parameters[name].to_dict()
|
|
for name in sorted(normalized_parameters)
|
|
},
|
|
"implementation_digest": _digest(implementation_digest, "$.implementation_digest"),
|
|
"input_schema_digest": normalized_input_digest,
|
|
"inputs": [item.to_dict() for item in normalized_inputs],
|
|
"valid_from": normalized_valid_from,
|
|
"valid_until": normalized_valid_until,
|
|
"warmup_sessions": normalized_warmup,
|
|
"lag_sessions": normalized_lag,
|
|
"producer": producer.to_dict(),
|
|
"code_revision": _git_revision(code_revision, "$.code_revision"),
|
|
}
|
|
expected_id = f"rhfactorv1:sha256:{hashlib.sha256(canonical_json_bytes(payload)).hexdigest()}"
|
|
if supplied_definition_id is not None:
|
|
normalized_supplied = _string(supplied_definition_id, "$.definition_id", _DEFINITION_ID)
|
|
if normalized_supplied != expected_id:
|
|
_fail(
|
|
ContractErrorCode.IDENTITY_MISMATCH,
|
|
"$.definition_id",
|
|
"does not identify the complete normalized definition",
|
|
)
|
|
instance = object.__new__(cls)
|
|
for name, value in {
|
|
**payload,
|
|
"definition_id": expected_id,
|
|
"parameters": MappingProxyType(dict(sorted(normalized_parameters.items()))),
|
|
"inputs": normalized_inputs,
|
|
"producer": producer,
|
|
}.items():
|
|
object.__setattr__(instance, name, value)
|
|
return instance
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"contract_name": self.contract_name,
|
|
"schema_version": self.schema_version,
|
|
"definition_id": self.definition_id,
|
|
"factor_id": self.factor_id,
|
|
"version": self.version,
|
|
"formula": self.formula,
|
|
"parameters": {name: value.to_dict() for name, value in self.parameters.items()},
|
|
"implementation_digest": self.implementation_digest,
|
|
"input_schema_digest": self.input_schema_digest,
|
|
"inputs": [item.to_dict() for item in self.inputs],
|
|
"valid_from": self.valid_from,
|
|
"valid_until": self.valid_until,
|
|
"warmup_sessions": self.warmup_sessions,
|
|
"lag_sessions": self.lag_sessions,
|
|
"producer": self.producer.to_dict(),
|
|
"code_revision": self.code_revision,
|
|
}
|
|
|
|
def to_json(self) -> str:
|
|
return canonical_json(self.to_dict())
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any) -> Self:
|
|
item = _object(
|
|
value,
|
|
"$",
|
|
(
|
|
"contract_name",
|
|
"schema_version",
|
|
"definition_id",
|
|
"factor_id",
|
|
"version",
|
|
"formula",
|
|
"parameters",
|
|
"implementation_digest",
|
|
"input_schema_digest",
|
|
"inputs",
|
|
"valid_from",
|
|
"valid_until",
|
|
"warmup_sessions",
|
|
"lag_sessions",
|
|
"producer",
|
|
"code_revision",
|
|
),
|
|
)
|
|
if item["contract_name"] != "researchhub.factor-definition":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.contract_name", "unsupported contract")
|
|
if item["schema_version"] != FACTOR_CONTRACT_VERSION:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.schema_version", "unsupported version")
|
|
if type(item["parameters"]) is not dict:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.parameters", "must be an object")
|
|
raw_parameters = item["parameters"]
|
|
parameters = {
|
|
name: TypedParameter.from_dict(parameter, f"$.parameters.{name}")
|
|
for name, parameter in raw_parameters.items()
|
|
}
|
|
raw_inputs = _array(item["inputs"], "$.inputs", minimum=1, unique=True)
|
|
return cls._build(
|
|
factor_id=item["factor_id"],
|
|
version=item["version"],
|
|
formula=item["formula"],
|
|
parameters=parameters,
|
|
implementation_digest=item["implementation_digest"],
|
|
input_schema_digest=item["input_schema_digest"],
|
|
inputs=tuple(
|
|
FactorInput.from_dict(raw_input, f"$.inputs[{index}]")
|
|
for index, raw_input in enumerate(raw_inputs)
|
|
),
|
|
valid_from=item["valid_from"],
|
|
valid_until=item["valid_until"],
|
|
warmup_sessions=item["warmup_sessions"],
|
|
lag_sessions=item["lag_sessions"],
|
|
producer=ProducerIdentity.from_dict(item["producer"]),
|
|
code_revision=item["code_revision"],
|
|
supplied_definition_id=item["definition_id"],
|
|
)
|
|
|
|
@classmethod
|
|
def from_json(cls, value: str | bytes) -> Self:
|
|
return cls.from_dict(_parse_json_object(value, "$"))
|
|
|
|
|
|
def validate_factor_catalog(definitions: Sequence[FactorDefinition]) -> tuple[FactorDefinition, ...]:
|
|
if type(definitions) not in {tuple, list}:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.definitions", "must be a sequence")
|
|
normalized: list[FactorDefinition] = []
|
|
for index, definition in enumerate(definitions):
|
|
if not isinstance(definition, FactorDefinition):
|
|
_fail(
|
|
ContractErrorCode.TYPE_ERROR,
|
|
f"$.definitions[{index}]",
|
|
"must be FactorDefinition",
|
|
)
|
|
normalized.append(definition)
|
|
if not normalized:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.definitions", "must not be empty")
|
|
identities = [definition.definition_id for definition in normalized]
|
|
if len(identities) != len(set(identities)):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.definitions", "duplicate definition_id")
|
|
by_logical: dict[tuple[str, str], list[FactorDefinition]] = defaultdict(list)
|
|
for definition in normalized:
|
|
by_logical[(definition.factor_id, definition.version)].append(definition)
|
|
for group in by_logical.values():
|
|
ordered = sorted(group, key=lambda definition: _parse_utc(definition.valid_from, "$.valid_from"))
|
|
for previous, current in pairwise(ordered):
|
|
if _parse_utc(current.valid_from, "$.valid_from") < _parse_utc(
|
|
previous.valid_until, "$.valid_until"
|
|
):
|
|
_fail(
|
|
ContractErrorCode.TIME_ORDER_VIOLATION,
|
|
"$.definitions",
|
|
"overlapping validity for the same factor_id/version",
|
|
)
|
|
return tuple(sorted(normalized, key=lambda definition: definition.definition_id))
|
|
|
|
|
|
def factor_definition_from_alpha158(
|
|
alpha_id: str,
|
|
*,
|
|
version: str,
|
|
parameters: Mapping[str, TypedParameter],
|
|
inputs: Sequence[FactorInput],
|
|
implementation_digest: str,
|
|
input_schema_digest: str,
|
|
valid_from: str,
|
|
valid_until: str,
|
|
warmup_sessions: int,
|
|
lag_sessions: int,
|
|
producer: ProducerIdentity,
|
|
code_revision: str,
|
|
) -> FactorDefinition:
|
|
"""Translate existing Alpha158 metadata without copying or executing a formula."""
|
|
|
|
if alpha_id not in ALPHA158_REGISTRY:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.alpha_id", "unknown Alpha158 registry key")
|
|
metadata = ALPHA158_REGISTRY[alpha_id]
|
|
declared_parameters = metadata["params"]
|
|
if declared_parameters:
|
|
expected_names = {
|
|
item if isinstance(item, str) else item.get("name")
|
|
for item in declared_parameters
|
|
}
|
|
if None in expected_names or set(parameters) != expected_names:
|
|
_fail(
|
|
ContractErrorCode.INPUT_CLOSURE_VIOLATION,
|
|
"$.parameters",
|
|
"does not exactly correspond to Alpha158 declared parameters",
|
|
)
|
|
elif parameters:
|
|
_fail(
|
|
ContractErrorCode.INPUT_CLOSURE_VIOLATION,
|
|
"$.parameters",
|
|
"registry declares no parameters",
|
|
)
|
|
normalized_inputs = _normalize_factor_inputs(inputs)
|
|
columns = [column for item in normalized_inputs for column in item.required_columns]
|
|
if len(columns) != len(set(columns)) or set(columns) != set(metadata["inputs"]):
|
|
_fail(
|
|
ContractErrorCode.INPUT_CLOSURE_VIOLATION,
|
|
"$.inputs",
|
|
"does not exactly correspond to Alpha158 declared input columns",
|
|
)
|
|
return FactorDefinition.create(
|
|
factor_id=alpha_id,
|
|
version=version,
|
|
formula=str(metadata["formula"]),
|
|
parameters=parameters,
|
|
implementation_digest=implementation_digest,
|
|
input_schema_digest=input_schema_digest,
|
|
inputs=normalized_inputs,
|
|
valid_from=valid_from,
|
|
valid_until=valid_until,
|
|
warmup_sessions=warmup_sessions,
|
|
lag_sessions=lag_sessions,
|
|
producer=producer,
|
|
code_revision=code_revision,
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class _SnapshotFacts:
|
|
snapshot_id: str
|
|
pit_cutoff: str
|
|
knowledge_start: datetime
|
|
knowledge_end: datetime
|
|
published_at: datetime
|
|
content_digest: str
|
|
manifest_digest: str
|
|
quality_status: str
|
|
quality_evidence_digests: tuple[str, ...]
|
|
qualification_status: str
|
|
qualification_policy_id: str
|
|
qualification_policy_version: str
|
|
qualification_evaluated_at: str
|
|
qualification_evidence_digest: str
|
|
|
|
|
|
def _validate_dataset_snapshot(document: dict[str, Any]) -> _SnapshotFacts:
|
|
_assert_canonical_profile(document)
|
|
root = _object(document, "$", ("contract_name", "schema_version", "snapshot_id", "descriptor"))
|
|
if root["contract_name"] != "researchhub.dataset-snapshot":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.contract_name", "unsupported upstream contract")
|
|
if root["schema_version"] != "1.0.0":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.schema_version", "unsupported upstream version")
|
|
snapshot_id = _string(root["snapshot_id"], "$.snapshot_id", _SNAPSHOT_ID)
|
|
descriptor = _object(
|
|
root["descriptor"],
|
|
"$.descriptor",
|
|
("dataset", "published_at", "time_semantics", "content", "lineage", "quality", "qualification"),
|
|
)
|
|
dataset = _object(
|
|
descriptor["dataset"],
|
|
"$.descriptor.dataset",
|
|
("dataset_id", "dataset_kind", "record_schema_version", "dimensions"),
|
|
)
|
|
dataset_id = _string(dataset["dataset_id"], "$.descriptor.dataset.dataset_id", _DATASET_ID)
|
|
dataset_kind = _enum(
|
|
dataset["dataset_kind"], "$.descriptor.dataset.dataset_kind", {"market", "macroeconomic"}
|
|
)
|
|
if not dataset_id.startswith(f"rhdataset:{dataset_kind}:"):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.dataset.dataset_id", "kind mismatch")
|
|
_semver(dataset["record_schema_version"], "$.descriptor.dataset.record_schema_version")
|
|
dimensions = _array(dataset["dimensions"], "$.descriptor.dataset.dimensions", minimum=1, unique=True)
|
|
expected_dimensions = (
|
|
["instrument_id", "effective_time"]
|
|
if dataset_kind == "market"
|
|
else ["series_id", "observation_period"]
|
|
)
|
|
if dimensions != expected_dimensions:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.dataset.dimensions", "kind dimensions mismatch")
|
|
for index, dimension in enumerate(dimensions):
|
|
_string(dimension, f"$.descriptor.dataset.dimensions[{index}]", _FIELD_NAME)
|
|
|
|
published_at = _parse_utc(descriptor["published_at"], "$.descriptor.published_at")
|
|
time_semantics = _object(
|
|
descriptor["time_semantics"],
|
|
"$.descriptor.time_semantics",
|
|
("effective_time", "knowledge_time", "pit_cutoff"),
|
|
)
|
|
ranges: dict[str, tuple[datetime, datetime]] = {}
|
|
for range_name in ("effective_time", "knowledge_time"):
|
|
range_value = _object(
|
|
time_semantics[range_name],
|
|
f"$.descriptor.time_semantics.{range_name}",
|
|
("start_inclusive", "end_inclusive"),
|
|
)
|
|
start = _parse_utc(
|
|
range_value["start_inclusive"],
|
|
f"$.descriptor.time_semantics.{range_name}.start_inclusive",
|
|
)
|
|
end = _parse_utc(
|
|
range_value["end_inclusive"],
|
|
f"$.descriptor.time_semantics.{range_name}.end_inclusive",
|
|
)
|
|
if start > end:
|
|
_fail(
|
|
ContractErrorCode.TIME_ORDER_VIOLATION,
|
|
f"$.descriptor.time_semantics.{range_name}",
|
|
"range must be ordered",
|
|
)
|
|
ranges[range_name] = (start, end)
|
|
pit_text = _string(time_semantics["pit_cutoff"], "$.descriptor.time_semantics.pit_cutoff", _UTC_INSTANT)
|
|
pit_cutoff = _parse_utc(pit_text, "$.descriptor.time_semantics.pit_cutoff")
|
|
|
|
content = _object(
|
|
descriptor["content"],
|
|
"$.descriptor.content",
|
|
(
|
|
"digest_algorithm",
|
|
"canonicalization",
|
|
"record_order",
|
|
"content_digest",
|
|
"logical_manifest",
|
|
"manifest_digest",
|
|
"record_count",
|
|
),
|
|
)
|
|
if content["digest_algorithm"] != "sha256" or content["canonicalization"] != "RFC8785":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.content", "unsupported digest profile")
|
|
if content["record_order"] != "canonical-record-byte-order":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.content.record_order", "unsupported order")
|
|
content_digest = _digest(content["content_digest"], "$.descriptor.content.content_digest")
|
|
record_count = _safe_integer(content["record_count"], "$.descriptor.content.record_count", minimum=0)
|
|
manifest = _object(
|
|
content["logical_manifest"],
|
|
"$.descriptor.content.logical_manifest",
|
|
("record_count", "chunks"),
|
|
)
|
|
manifest_count = _safe_integer(
|
|
manifest["record_count"], "$.descriptor.content.logical_manifest.record_count", minimum=0
|
|
)
|
|
chunks = _array(
|
|
manifest["chunks"], "$.descriptor.content.logical_manifest.chunks", minimum=1, unique=True
|
|
)
|
|
chunk_count = 0
|
|
for index, raw_chunk in enumerate(chunks):
|
|
chunk = _object(
|
|
raw_chunk,
|
|
f"$.descriptor.content.logical_manifest.chunks[{index}]",
|
|
("chunk_index", "content_digest", "record_count"),
|
|
)
|
|
if _safe_integer(chunk["chunk_index"], f"$.descriptor.content.logical_manifest.chunks[{index}].chunk_index", minimum=0) != index:
|
|
_fail(
|
|
ContractErrorCode.INVALID_VALUE,
|
|
f"$.descriptor.content.logical_manifest.chunks[{index}].chunk_index",
|
|
"chunk indexes must be contiguous",
|
|
)
|
|
_digest(chunk["content_digest"], f"$.descriptor.content.logical_manifest.chunks[{index}].content_digest")
|
|
chunk_count += _safe_integer(
|
|
chunk["record_count"],
|
|
f"$.descriptor.content.logical_manifest.chunks[{index}].record_count",
|
|
minimum=0,
|
|
)
|
|
if manifest_count != record_count or chunk_count != record_count:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.content.logical_manifest", "record counts mismatch")
|
|
manifest_digest = _digest(content["manifest_digest"], "$.descriptor.content.manifest_digest")
|
|
expected_manifest_digest = _digest_bytes(canonical_json_bytes(manifest))
|
|
if manifest_digest != expected_manifest_digest:
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, "$.descriptor.content.manifest_digest", "manifest mismatch")
|
|
|
|
lineage = _object(
|
|
descriptor["lineage"],
|
|
"$.descriptor.lineage",
|
|
("publisher", "transformation", "upstream_snapshot_ids", "upstream_content_digests"),
|
|
)
|
|
publisher = _object(lineage["publisher"], "$.descriptor.lineage.publisher", ("id", "version"))
|
|
if publisher["id"] != "researchhub.data":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.lineage.publisher.id", "wrong authority")
|
|
_semver(publisher["version"], "$.descriptor.lineage.publisher.version")
|
|
transformation = _object(
|
|
lineage["transformation"], "$.descriptor.lineage.transformation", ("id", "version")
|
|
)
|
|
_string(transformation["id"], "$.descriptor.lineage.transformation.id", _TRANSFORMATION_ID)
|
|
_semver(transformation["version"], "$.descriptor.lineage.transformation.version")
|
|
for key, pattern in (("upstream_snapshot_ids", _SNAPSHOT_ID), ("upstream_content_digests", _SHA256)):
|
|
values = _array(lineage[key], f"$.descriptor.lineage.{key}", unique=True)
|
|
for index, value in enumerate(values):
|
|
_string(value, f"$.descriptor.lineage.{key}[{index}]", pattern)
|
|
|
|
quality = _object(descriptor["quality"], "$.descriptor.quality", ("status", "checks"))
|
|
quality_status = _enum(quality["status"], "$.descriptor.quality.status", {"passed", "failed"})
|
|
checks = _array(quality["checks"], "$.descriptor.quality.checks", minimum=1)
|
|
check_ids: set[str] = set()
|
|
quality_evidence: list[str] = []
|
|
all_checks_passed = True
|
|
for index, raw_check in enumerate(checks):
|
|
check = _object(
|
|
raw_check,
|
|
f"$.descriptor.quality.checks[{index}]",
|
|
("check_id", "status", "severity", "evidence_digest"),
|
|
)
|
|
check_id = _enum(
|
|
check["check_id"],
|
|
f"$.descriptor.quality.checks[{index}].check_id",
|
|
{"completeness", "duplicate_identity", "pit_time_integrity", "range_validity", "schema_conformance"},
|
|
)
|
|
if check_id in check_ids:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.quality.checks", "duplicate check_id")
|
|
check_ids.add(check_id)
|
|
check_status = _enum(
|
|
check["status"], f"$.descriptor.quality.checks[{index}].status", {"passed", "failed"}
|
|
)
|
|
all_checks_passed = all_checks_passed and check_status == "passed"
|
|
_enum(
|
|
check["severity"],
|
|
f"$.descriptor.quality.checks[{index}].severity",
|
|
{"blocking", "advisory"},
|
|
)
|
|
quality_evidence.append(
|
|
_digest(check["evidence_digest"], f"$.descriptor.quality.checks[{index}].evidence_digest")
|
|
)
|
|
|
|
qualification = _object(
|
|
descriptor["qualification"],
|
|
"$.descriptor.qualification",
|
|
("status", "policy_id", "policy_version", "evaluated_at", "evidence_digest"),
|
|
)
|
|
qualification_status = _enum(
|
|
qualification["status"], "$.descriptor.qualification.status", {"qualified", "rejected"}
|
|
)
|
|
if qualification["policy_id"] != "researchhub.dataset-snapshot.pit":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.descriptor.qualification.policy_id", "wrong policy")
|
|
policy_version = _semver(qualification["policy_version"], "$.descriptor.qualification.policy_version")
|
|
evaluated_text = _string(
|
|
qualification["evaluated_at"], "$.descriptor.qualification.evaluated_at", _UTC_INSTANT
|
|
)
|
|
evaluated_at = _parse_utc(evaluated_text, "$.descriptor.qualification.evaluated_at")
|
|
qualification_evidence = _digest(
|
|
qualification["evidence_digest"], "$.descriptor.qualification.evidence_digest"
|
|
)
|
|
knowledge_start, knowledge_end = ranges["knowledge_time"]
|
|
if not knowledge_start <= knowledge_end <= pit_cutoff <= evaluated_at <= published_at:
|
|
_fail(
|
|
ContractErrorCode.TIME_ORDER_VIOLATION,
|
|
"$.descriptor.time_semantics",
|
|
"knowledge <= PIT <= qualification <= publication is required",
|
|
)
|
|
if qualification_status == "qualified" and (
|
|
quality_status != "passed" or not all_checks_passed
|
|
):
|
|
_fail(
|
|
ContractErrorCode.QUALIFICATION_REJECTED,
|
|
"$.descriptor.qualification",
|
|
"qualified snapshots require all quality checks passed",
|
|
)
|
|
if quality_status == "failed" and qualification_status != "rejected":
|
|
_fail(ContractErrorCode.QUALIFICATION_REJECTED, "$.descriptor.qualification", "failed quality")
|
|
expected_snapshot_id = _content_address(document, "snapshot_id", "rhdsv1:sha256:")
|
|
if snapshot_id != expected_snapshot_id:
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, "$.snapshot_id", "snapshot envelope mismatch")
|
|
return _SnapshotFacts(
|
|
snapshot_id=snapshot_id,
|
|
pit_cutoff=pit_text,
|
|
knowledge_start=knowledge_start,
|
|
knowledge_end=knowledge_end,
|
|
published_at=published_at,
|
|
content_digest=content_digest,
|
|
manifest_digest=manifest_digest,
|
|
quality_status=quality_status,
|
|
quality_evidence_digests=tuple(quality_evidence),
|
|
qualification_status=qualification_status,
|
|
qualification_policy_id="researchhub.dataset-snapshot.pit",
|
|
qualification_policy_version=policy_version,
|
|
qualification_evaluated_at=evaluated_text,
|
|
qualification_evidence_digest=qualification_evidence,
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True, init=False)
|
|
class DatasetSnapshotEnvelope:
|
|
snapshot_id: str
|
|
pit_cutoff: str
|
|
knowledge_start: datetime
|
|
knowledge_end: datetime
|
|
published_at: datetime
|
|
content_digest: str
|
|
manifest_digest: str
|
|
quality_status: str
|
|
quality_evidence_digests: tuple[str, ...]
|
|
qualification_status: str
|
|
qualification_policy_id: str
|
|
qualification_policy_version: str
|
|
qualification_evaluated_at: str
|
|
qualification_evidence_digest: str
|
|
_payload: Mapping[str, Any] = field(repr=False, compare=False)
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any) -> Self:
|
|
if type(value) is not dict:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$", "DatasetSnapshot envelope must be an object")
|
|
detached = _thaw_json(_freeze_json(value))
|
|
facts = _validate_dataset_snapshot(detached)
|
|
instance = object.__new__(cls)
|
|
for fact_field in fields(facts):
|
|
object.__setattr__(instance, fact_field.name, getattr(facts, fact_field.name))
|
|
object.__setattr__(instance, "_payload", _freeze_json(detached))
|
|
return instance
|
|
|
|
@classmethod
|
|
def from_json(cls, value: str | bytes) -> Self:
|
|
return cls.from_dict(_parse_json_object(value, "$"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return cast(dict[str, Any], _thaw_json(self._payload))
|
|
|
|
def to_json(self) -> str:
|
|
return canonical_json(self.to_dict())
|
|
|
|
def require_qualified(self) -> None:
|
|
if self.qualification_status != "qualified" or self.quality_status != "passed":
|
|
_fail(
|
|
ContractErrorCode.QUALIFICATION_REJECTED,
|
|
"$.dataset_snapshot.descriptor.qualification",
|
|
"only jointly qualified, passed snapshots are admissible",
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class _ViewFacts:
|
|
view_ref_id: str
|
|
schema_digest: str
|
|
content_digest: str
|
|
transformation_digest: str
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class _FoundationFacts:
|
|
foundation_id: str
|
|
dataset_snapshot_id: str
|
|
pit_cutoff: str
|
|
views: Mapping[str, _ViewFacts]
|
|
evidence_scope: str
|
|
contract_evidence_digests: tuple[str, ...]
|
|
real_data_status: str
|
|
|
|
|
|
def _validate_revision_chain(
|
|
records: list[dict[str, Any]],
|
|
*,
|
|
identity_field: str,
|
|
group_fields: tuple[str, ...],
|
|
path: str,
|
|
) -> None:
|
|
groups: dict[tuple[Any, ...], list[dict[str, Any]]] = defaultdict(list)
|
|
for record in records:
|
|
groups[tuple(record[name] for name in group_fields)].append(record)
|
|
for group in groups.values():
|
|
ordered = sorted(group, key=lambda record: record["revision_number"])
|
|
if [record["revision_number"] for record in ordered] != list(range(1, len(ordered) + 1)):
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, path, "revision chain contains a gap")
|
|
if "supersedes_revision_id" in ordered[0]:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, path, "first revision has a parent")
|
|
for previous, current in pairwise(ordered):
|
|
if current.get("supersedes_revision_id") != previous[identity_field]:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, path, "revision ancestry is not contiguous")
|
|
if _parse_utc(current["knowledge_time"], path) <= _parse_utc(
|
|
previous["knowledge_time"], path
|
|
):
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, path, "revision knowledge must increase")
|
|
|
|
|
|
def _validate_data_foundation(document: dict[str, Any]) -> _FoundationFacts:
|
|
_assert_canonical_profile(document)
|
|
_no_physical_leakage(document)
|
|
root = _object(
|
|
document,
|
|
"$",
|
|
(
|
|
"contract_name",
|
|
"schema_version",
|
|
"foundation_id",
|
|
"dataset_snapshot_id",
|
|
"pit_cutoff",
|
|
"instrument_routes",
|
|
"trading_calendar_revisions",
|
|
"corporate_action_revisions",
|
|
"standardized_views",
|
|
"revision_lineage",
|
|
"readiness",
|
|
),
|
|
)
|
|
if root["contract_name"] != "researchhub.data-foundation":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.contract_name", "unsupported upstream contract")
|
|
if root["schema_version"] != "1.0.0":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.schema_version", "unsupported upstream version")
|
|
foundation_id = _string(root["foundation_id"], "$.foundation_id", _FOUNDATION_ID)
|
|
snapshot_id = _string(root["dataset_snapshot_id"], "$.dataset_snapshot_id", _SNAPSHOT_ID)
|
|
pit_text = _string(root["pit_cutoff"], "$.pit_cutoff", _UTC_INSTANT)
|
|
pit_cutoff = _parse_utc(pit_text, "$.pit_cutoff")
|
|
|
|
routes_raw = _array(root["instrument_routes"], "$.instrument_routes", minimum=1, unique=True)
|
|
routes: list[dict[str, Any]] = []
|
|
for index, raw_route in enumerate(routes_raw):
|
|
path = f"$.instrument_routes[{index}]"
|
|
route = _object(
|
|
raw_route,
|
|
path,
|
|
(
|
|
"route_revision_id",
|
|
"instrument_id",
|
|
"revision_number",
|
|
"symbol",
|
|
"mic",
|
|
"currency",
|
|
"asset_class",
|
|
"instrument_type",
|
|
"calendar_id",
|
|
"effective_from",
|
|
"knowledge_time",
|
|
"evidence_digest",
|
|
),
|
|
("supersedes_revision_id",),
|
|
)
|
|
_string(route["route_revision_id"], f"{path}.route_revision_id", _ROUTE_ID)
|
|
_string(route["instrument_id"], f"{path}.instrument_id", _INSTRUMENT_ID)
|
|
_safe_integer(route["revision_number"], f"{path}.revision_number", minimum=1)
|
|
if "supersedes_revision_id" in route:
|
|
_string(route["supersedes_revision_id"], f"{path}.supersedes_revision_id", _ROUTE_ID)
|
|
symbol = _string(route["symbol"], f"{path}.symbol")
|
|
if not 1 <= len(symbol) <= 32 or re.fullmatch(r"^[A-Z0-9][A-Z0-9.-]*$", symbol) is None:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, f"{path}.symbol", "invalid symbol profile")
|
|
if re.fullmatch(r"^[A-Z0-9]{4}$", _string(route["mic"], f"{path}.mic")) is None:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, f"{path}.mic", "invalid MIC")
|
|
if re.fullmatch(r"^[A-Z]{3}$", _string(route["currency"], f"{path}.currency")) is None:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, f"{path}.currency", "invalid currency")
|
|
_enum(route["asset_class"], f"{path}.asset_class", {"equity", "fund", "fixed_income", "future", "option"})
|
|
_enum(route["instrument_type"], f"{path}.instrument_type", {"stock", "etf", "bond", "future", "option"})
|
|
_string(route["calendar_id"], f"{path}.calendar_id", _CALENDAR_ID)
|
|
_parse_utc(route["effective_from"], f"{path}.effective_from")
|
|
_parse_utc(route["knowledge_time"], f"{path}.knowledge_time")
|
|
_digest(route["evidence_digest"], f"{path}.evidence_digest")
|
|
if route["route_revision_id"] != _content_address(route, "route_revision_id", "rhroutev1:sha256:"):
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, f"{path}.route_revision_id", "route mismatch")
|
|
routes.append(route)
|
|
|
|
calendars_raw = _array(
|
|
root["trading_calendar_revisions"], "$.trading_calendar_revisions", minimum=1, unique=True
|
|
)
|
|
calendars: list[dict[str, Any]] = []
|
|
for index, raw_calendar in enumerate(calendars_raw):
|
|
path = f"$.trading_calendar_revisions[{index}]"
|
|
calendar = _object(
|
|
raw_calendar,
|
|
path,
|
|
(
|
|
"calendar_revision_id",
|
|
"calendar_id",
|
|
"session_date",
|
|
"revision_number",
|
|
"status",
|
|
"sessions",
|
|
"knowledge_time",
|
|
"evidence_digest",
|
|
),
|
|
("supersedes_revision_id",),
|
|
)
|
|
_string(calendar["calendar_revision_id"], f"{path}.calendar_revision_id", _CALENDAR_REVISION_ID)
|
|
_string(calendar["calendar_id"], f"{path}.calendar_id", _CALENDAR_ID)
|
|
_parse_date(calendar["session_date"], f"{path}.session_date")
|
|
_safe_integer(calendar["revision_number"], f"{path}.revision_number", minimum=1)
|
|
if "supersedes_revision_id" in calendar:
|
|
_string(calendar["supersedes_revision_id"], f"{path}.supersedes_revision_id", _CALENDAR_REVISION_ID)
|
|
status = _enum(calendar["status"], f"{path}.status", {"open", "closed"})
|
|
sessions = _array(calendar["sessions"], f"{path}.sessions", unique=True)
|
|
if status == "open" and not sessions:
|
|
_fail(ContractErrorCode.INVALID_VALUE, f"{path}.sessions", "open session requires segments")
|
|
if status == "closed" and sessions:
|
|
_fail(ContractErrorCode.INVALID_VALUE, f"{path}.sessions", "closed session must be empty")
|
|
parsed_sessions: list[tuple[datetime, datetime]] = []
|
|
for session_index, raw_session in enumerate(sessions):
|
|
session_path = f"{path}.sessions[{session_index}]"
|
|
session = _object(raw_session, session_path, ("opens_at", "closes_at"))
|
|
opens_at = _parse_utc(session["opens_at"], f"{session_path}.opens_at")
|
|
closes_at = _parse_utc(session["closes_at"], f"{session_path}.closes_at")
|
|
if closes_at <= opens_at:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, session_path, "session close must follow open")
|
|
parsed_sessions.append((opens_at, closes_at))
|
|
if parsed_sessions != sorted(parsed_sessions) or any(
|
|
current[0] < previous[1]
|
|
for previous, current in pairwise(parsed_sessions)
|
|
):
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, f"{path}.sessions", "sessions overlap or are unordered")
|
|
_parse_utc(calendar["knowledge_time"], f"{path}.knowledge_time")
|
|
_digest(calendar["evidence_digest"], f"{path}.evidence_digest")
|
|
if calendar["calendar_revision_id"] != _content_address(
|
|
calendar, "calendar_revision_id", "rhcalv1:sha256:"
|
|
):
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, f"{path}.calendar_revision_id", "calendar mismatch")
|
|
calendars.append(calendar)
|
|
|
|
actions_raw = _array(root["corporate_action_revisions"], "$.corporate_action_revisions", unique=True)
|
|
actions: list[dict[str, Any]] = []
|
|
for index, raw_action in enumerate(actions_raw):
|
|
path = f"$.corporate_action_revisions[{index}]"
|
|
action = _object(
|
|
raw_action,
|
|
path,
|
|
(
|
|
"action_revision_id",
|
|
"action_id",
|
|
"instrument_id",
|
|
"revision_number",
|
|
"action_type",
|
|
"status",
|
|
"effective_time",
|
|
"knowledge_time",
|
|
"terms_digest",
|
|
"evidence_digest",
|
|
),
|
|
("supersedes_revision_id",),
|
|
)
|
|
_string(action["action_revision_id"], f"{path}.action_revision_id", _ACTION_REVISION_ID)
|
|
_string(action["action_id"], f"{path}.action_id", _ACTION_ID)
|
|
_string(action["instrument_id"], f"{path}.instrument_id", _INSTRUMENT_ID)
|
|
_safe_integer(action["revision_number"], f"{path}.revision_number", minimum=1)
|
|
if "supersedes_revision_id" in action:
|
|
_string(action["supersedes_revision_id"], f"{path}.supersedes_revision_id", _ACTION_REVISION_ID)
|
|
_enum(action["action_type"], f"{path}.action_type", {"cash_dividend", "stock_dividend", "split", "rights_issue", "symbol_change", "delisting"})
|
|
_enum(action["status"], f"{path}.status", {"announced", "confirmed", "cancelled"})
|
|
_parse_utc(action["effective_time"], f"{path}.effective_time")
|
|
_parse_utc(action["knowledge_time"], f"{path}.knowledge_time")
|
|
_digest(action["terms_digest"], f"{path}.terms_digest")
|
|
_digest(action["evidence_digest"], f"{path}.evidence_digest")
|
|
if action["action_revision_id"] != _content_address(
|
|
action, "action_revision_id", "rhcav1:sha256:"
|
|
):
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, f"{path}.action_revision_id", "action mismatch")
|
|
actions.append(action)
|
|
|
|
_validate_revision_chain(routes, identity_field="route_revision_id", group_fields=("instrument_id",), path="$.instrument_routes")
|
|
_validate_revision_chain(calendars, identity_field="calendar_revision_id", group_fields=("calendar_id", "session_date"), path="$.trading_calendar_revisions")
|
|
_validate_revision_chain(actions, identity_field="action_revision_id", group_fields=("action_id",), path="$.corporate_action_revisions")
|
|
|
|
route_by_id = {route["route_revision_id"]: route for route in routes}
|
|
calendar_by_id = {calendar["calendar_revision_id"]: calendar for calendar in calendars}
|
|
action_by_id = {action["action_revision_id"]: action for action in actions}
|
|
if len(route_by_id) != len(routes) or len(calendar_by_id) != len(calendars) or len(action_by_id) != len(actions):
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$", "duplicate revision identity")
|
|
calendar_ids = {calendar["calendar_id"] for calendar in calendars}
|
|
instruments = {route["instrument_id"] for route in routes}
|
|
if any(route["calendar_id"] not in calendar_ids for route in routes):
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.instrument_routes", "route calendar missing")
|
|
if any(action["instrument_id"] not in instruments for action in actions):
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.corporate_action_revisions", "instrument route missing")
|
|
|
|
lineage_raw = _array(root["revision_lineage"], "$.revision_lineage", minimum=1, unique=True)
|
|
lineage: dict[str, dict[str, Any]] = {}
|
|
allowed_revision_ids = set(route_by_id) | set(calendar_by_id) | set(action_by_id)
|
|
for index, raw_entry in enumerate(lineage_raw):
|
|
path = f"$.revision_lineage[{index}]"
|
|
entry = _object(
|
|
raw_entry,
|
|
path,
|
|
("revision_kind", "revision_id", "revision_number", "knowledge_time", "evidence_digest"),
|
|
("supersedes_revision_id",),
|
|
)
|
|
kind = _enum(entry["revision_kind"], f"{path}.revision_kind", {"instrument_route", "trading_calendar", "corporate_action"})
|
|
pattern = {
|
|
"instrument_route": _ROUTE_ID,
|
|
"trading_calendar": _CALENDAR_REVISION_ID,
|
|
"corporate_action": _ACTION_REVISION_ID,
|
|
}[kind]
|
|
revision_id = _string(entry["revision_id"], f"{path}.revision_id", pattern)
|
|
if "supersedes_revision_id" in entry:
|
|
_string(entry["supersedes_revision_id"], f"{path}.supersedes_revision_id", pattern)
|
|
_safe_integer(entry["revision_number"], f"{path}.revision_number", minimum=1)
|
|
_parse_utc(entry["knowledge_time"], f"{path}.knowledge_time")
|
|
_digest(entry["evidence_digest"], f"{path}.evidence_digest")
|
|
if revision_id in lineage:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, path, "duplicate lineage identity")
|
|
lineage[revision_id] = entry
|
|
if set(lineage) != allowed_revision_ids:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.revision_lineage", "must cover every revision exactly")
|
|
for kind, records, identity_field in (
|
|
("instrument_route", routes, "route_revision_id"),
|
|
("trading_calendar", calendars, "calendar_revision_id"),
|
|
("corporate_action", actions, "action_revision_id"),
|
|
):
|
|
for revision in records:
|
|
entry = lineage[revision[identity_field]]
|
|
if entry["revision_kind"] != kind or any(
|
|
entry[name] != revision[name]
|
|
for name in ("revision_number", "knowledge_time", "evidence_digest")
|
|
) or entry.get("supersedes_revision_id") != revision.get("supersedes_revision_id"):
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.revision_lineage", "entry mismatch")
|
|
if _parse_utc(revision["knowledge_time"], "$.revision_lineage.knowledge_time") > pit_cutoff:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.revision_lineage.knowledge_time", "future knowledge exceeds PIT")
|
|
|
|
views_raw = _array(root["standardized_views"], "$.standardized_views", minimum=1, unique=True)
|
|
views: dict[str, _ViewFacts] = {}
|
|
for index, raw_view in enumerate(views_raw):
|
|
path = f"$.standardized_views[{index}]"
|
|
view = _object(
|
|
raw_view,
|
|
path,
|
|
(
|
|
"view_ref_id",
|
|
"view_id",
|
|
"view_version",
|
|
"dataset_snapshot_id",
|
|
"pit_cutoff",
|
|
"schema_digest",
|
|
"content_digest",
|
|
"transformation_digest",
|
|
"instrument_route_revision_ids",
|
|
"trading_calendar_revision_ids",
|
|
"corporate_action_revision_ids",
|
|
),
|
|
)
|
|
view_ref_id = _string(view["view_ref_id"], f"{path}.view_ref_id", _VIEW_REF_ID)
|
|
_string(view["view_id"], f"{path}.view_id", _VIEW_ID)
|
|
_semver(view["view_version"], f"{path}.view_version")
|
|
if _string(view["dataset_snapshot_id"], f"{path}.dataset_snapshot_id", _SNAPSHOT_ID) != snapshot_id:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, f"{path}.dataset_snapshot_id", "view is not snapshot-bound")
|
|
if _string(view["pit_cutoff"], f"{path}.pit_cutoff", _UTC_INSTANT) != pit_text:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, f"{path}.pit_cutoff", "view uses a different PIT")
|
|
schema_digest = _digest(view["schema_digest"], f"{path}.schema_digest")
|
|
content_digest = _digest(view["content_digest"], f"{path}.content_digest")
|
|
transformation_digest = _digest(view["transformation_digest"], f"{path}.transformation_digest")
|
|
selected: list[tuple[list[Any], Mapping[str, dict[str, Any]], re.Pattern[str], str]] = [
|
|
(_array(view["instrument_route_revision_ids"], f"{path}.instrument_route_revision_ids", minimum=1, unique=True), route_by_id, _ROUTE_ID, "instrument_route_revision_ids"),
|
|
(_array(view["trading_calendar_revision_ids"], f"{path}.trading_calendar_revision_ids", minimum=1, unique=True), calendar_by_id, _CALENDAR_REVISION_ID, "trading_calendar_revision_ids"),
|
|
(_array(view["corporate_action_revision_ids"], f"{path}.corporate_action_revision_ids", unique=True), action_by_id, _ACTION_REVISION_ID, "corporate_action_revision_ids"),
|
|
]
|
|
for ids, known, pattern, field_name in selected:
|
|
for selected_index, revision_id in enumerate(ids):
|
|
normalized_revision_id = _string(revision_id, f"{path}.{field_name}[{selected_index}]", pattern)
|
|
if normalized_revision_id not in known:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, f"{path}.{field_name}", "unknown revision")
|
|
parent = lineage[normalized_revision_id].get("supersedes_revision_id")
|
|
while parent is not None:
|
|
if parent not in ids:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, f"{path}.{field_name}", "revision ancestry omitted")
|
|
parent = lineage[parent].get("supersedes_revision_id")
|
|
route_calendar_ids = {
|
|
route_by_id[revision_id]["calendar_id"]
|
|
for revision_id in view["instrument_route_revision_ids"]
|
|
}
|
|
selected_calendar_ids = {
|
|
calendar_by_id[revision_id]["calendar_id"]
|
|
for revision_id in view["trading_calendar_revision_ids"]
|
|
}
|
|
if not route_calendar_ids.issubset(selected_calendar_ids):
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, path, "selected route calendar is not selected by view")
|
|
if view_ref_id != _content_address(view, "view_ref_id", "rhviewrefv1:sha256:"):
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, f"{path}.view_ref_id", "view mismatch")
|
|
if view_ref_id in views:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.standardized_views", "duplicate view_ref_id")
|
|
views[view_ref_id] = _ViewFacts(view_ref_id, schema_digest, content_digest, transformation_digest)
|
|
|
|
readiness = _object(
|
|
root["readiness"],
|
|
"$.readiness",
|
|
("evidence_scope", "contract_validation", "real_data_validation", "production_validation", "live_validation"),
|
|
)
|
|
evidence_scope = _enum(readiness["evidence_scope"], "$.readiness.evidence_scope", {"synthetic_fixture", "real_data"})
|
|
contract_level = _object(readiness["contract_validation"], "$.readiness.contract_validation", ("status", "evidence_digests"))
|
|
if contract_level["status"] != "validated":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.readiness.contract_validation.status", "must be validated")
|
|
contract_evidence_raw = _array(contract_level["evidence_digests"], "$.readiness.contract_validation.evidence_digests", minimum=1, unique=True)
|
|
contract_evidence = tuple(_digest(value, f"$.readiness.contract_validation.evidence_digests[{index}]") for index, value in enumerate(contract_evidence_raw))
|
|
used_evidence = set(contract_evidence)
|
|
level_status: dict[str, str] = {}
|
|
for name in ("real_data_validation", "production_validation", "live_validation"):
|
|
level = _object(readiness[name], f"$.readiness.{name}", ("status", "evidence_digests"))
|
|
status = _enum(level["status"], f"$.readiness.{name}.status", {"not_validated", "validated"})
|
|
evidence_raw = _array(level["evidence_digests"], f"$.readiness.{name}.evidence_digests", unique=True)
|
|
evidence = {_digest(value, f"$.readiness.{name}.evidence_digests[{index}]") for index, value in enumerate(evidence_raw)}
|
|
if status == "validated" and not evidence:
|
|
_fail(ContractErrorCode.INVALID_VALUE, f"$.readiness.{name}", "validated level requires evidence")
|
|
if status == "not_validated" and evidence:
|
|
_fail(ContractErrorCode.INVALID_VALUE, f"$.readiness.{name}", "unvalidated level cannot have evidence")
|
|
if used_evidence.intersection(evidence):
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, f"$.readiness.{name}", "evidence reused across levels")
|
|
used_evidence.update(evidence)
|
|
level_status[name] = status
|
|
if evidence_scope == "synthetic_fixture" and any(status == "validated" for status in level_status.values()):
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, "$.readiness", "synthetic evidence cannot promote readiness")
|
|
if level_status["production_validation"] == "validated" and level_status["real_data_validation"] != "validated":
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, "$.readiness.production_validation", "real-data prerequisite missing")
|
|
if level_status["live_validation"] == "validated" and level_status["production_validation"] != "validated":
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, "$.readiness.live_validation", "production prerequisite missing")
|
|
if foundation_id != _content_address(document, "foundation_id", "rhdfv1:sha256:"):
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, "$.foundation_id", "foundation envelope mismatch")
|
|
return _FoundationFacts(
|
|
foundation_id=foundation_id,
|
|
dataset_snapshot_id=snapshot_id,
|
|
pit_cutoff=pit_text,
|
|
views=MappingProxyType(dict(views)),
|
|
evidence_scope=evidence_scope,
|
|
contract_evidence_digests=contract_evidence,
|
|
real_data_status=level_status["real_data_validation"],
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True, init=False)
|
|
class DataFoundationEnvelope:
|
|
foundation_id: str
|
|
dataset_snapshot_id: str
|
|
pit_cutoff: str
|
|
views: Mapping[str, _ViewFacts]
|
|
evidence_scope: str
|
|
contract_evidence_digests: tuple[str, ...]
|
|
real_data_status: str
|
|
_payload: Mapping[str, Any] = field(repr=False, compare=False)
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any) -> Self:
|
|
if type(value) is not dict:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$", "Data Foundation envelope must be an object")
|
|
detached = _thaw_json(_freeze_json(value))
|
|
facts = _validate_data_foundation(detached)
|
|
instance = object.__new__(cls)
|
|
for fact_field in fields(facts):
|
|
object.__setattr__(instance, fact_field.name, getattr(facts, fact_field.name))
|
|
object.__setattr__(instance, "_payload", _freeze_json(detached))
|
|
return instance
|
|
|
|
@classmethod
|
|
def from_json(cls, value: str | bytes) -> Self:
|
|
return cls.from_dict(_parse_json_object(value, "$"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return cast(dict[str, Any], _thaw_json(self._payload))
|
|
|
|
def to_json(self) -> str:
|
|
return canonical_json(self.to_dict())
|
|
|
|
|
|
class AvailabilityMode(StrEnum):
|
|
AS_AVAILABLE = "as_available"
|
|
RETROSPECTIVE_REPLAY = "retrospective_replay"
|
|
|
|
|
|
class HistoricalAvailability(StrEnum):
|
|
DECLARED_AS_AVAILABLE = "declared_as_available"
|
|
NOT_ESTABLISHED = "not_established"
|
|
|
|
|
|
class PayloadValidation(StrEnum):
|
|
PAYLOAD_REVALIDATED = "payload_revalidated"
|
|
REFERENCE_ONLY = "reference_only"
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ActorIdentity:
|
|
kind: str
|
|
id: str
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "kind", _enum(self.kind, "$.actor.kind", {"user", "service"}))
|
|
object.__setattr__(self, "id", _logical_id(self.id, "$.actor.id"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {"kind": self.kind, "id": self.id}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.actor") -> Self:
|
|
item = _object(value, path, ("kind", "id"))
|
|
return cls(kind=_string(item["kind"], f"{path}.kind"), id=_string(item["id"], f"{path}.id"))
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class Causation:
|
|
kind: str
|
|
id: str
|
|
|
|
def __post_init__(self) -> None:
|
|
kind = _enum(self.kind, "$.causation.kind", {"factor_set", "foundation"})
|
|
pattern = _FACTOR_SET_ID if kind == "factor_set" else _FOUNDATION_ID
|
|
object.__setattr__(self, "kind", kind)
|
|
object.__setattr__(self, "id", _string(self.id, "$.causation.id", pattern))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {"kind": self.kind, "id": self.id}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.causation") -> Self:
|
|
item = _object(value, path, ("kind", "id"))
|
|
return cls(kind=_string(item["kind"], f"{path}.kind"), id=_string(item["id"], f"{path}.id"))
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class InputBinding:
|
|
definition_id: str
|
|
input_name: str
|
|
view_ref_id: str
|
|
schema_digest: str
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "definition_id", _string(self.definition_id, "$.input_bindings[].definition_id", _DEFINITION_ID))
|
|
object.__setattr__(self, "input_name", _string(self.input_name, "$.input_bindings[].input_name", _FIELD_NAME))
|
|
object.__setattr__(self, "view_ref_id", _string(self.view_ref_id, "$.input_bindings[].view_ref_id", _VIEW_REF_ID))
|
|
object.__setattr__(self, "schema_digest", _digest(self.schema_digest, "$.input_bindings[].schema_digest"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"definition_id": self.definition_id,
|
|
"input_name": self.input_name,
|
|
"view_ref_id": self.view_ref_id,
|
|
"schema_digest": self.schema_digest,
|
|
}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str) -> Self:
|
|
item = _object(value, path, ("definition_id", "input_name", "view_ref_id", "schema_digest"))
|
|
return cls(
|
|
definition_id=_string(item["definition_id"], f"{path}.definition_id"),
|
|
input_name=_string(item["input_name"], f"{path}.input_name"),
|
|
view_ref_id=_string(item["view_ref_id"], f"{path}.view_ref_id"),
|
|
schema_digest=_string(item["schema_digest"], f"{path}.schema_digest"),
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ViewAvailability:
|
|
view_ref_id: str
|
|
available_at: str
|
|
evidence_digest: str
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "view_ref_id", _string(self.view_ref_id, "$.view_availability[].view_ref_id", _VIEW_REF_ID))
|
|
object.__setattr__(self, "available_at", _normalize_new_instant(self.available_at, "$.view_availability[].available_at"))
|
|
object.__setattr__(self, "evidence_digest", _digest(self.evidence_digest, "$.view_availability[].evidence_digest"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"view_ref_id": self.view_ref_id,
|
|
"available_at": self.available_at,
|
|
"evidence_digest": self.evidence_digest,
|
|
}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str) -> Self:
|
|
item = _object(value, path, ("view_ref_id", "available_at", "evidence_digest"))
|
|
return cls(
|
|
view_ref_id=_string(item["view_ref_id"], f"{path}.view_ref_id"),
|
|
available_at=_string(item["available_at"], f"{path}.available_at"),
|
|
evidence_digest=_string(item["evidence_digest"], f"{path}.evidence_digest"),
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class OutputQualityCheck:
|
|
check_id: str
|
|
status: str
|
|
evidence_digest: str
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "check_id", _string(self.check_id, "$.output_quality.checks[].check_id", _FIELD_NAME))
|
|
object.__setattr__(self, "status", _enum(self.status, "$.output_quality.checks[].status", {"failed", "passed"}))
|
|
object.__setattr__(self, "evidence_digest", _digest(self.evidence_digest, "$.output_quality.checks[].evidence_digest"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {"check_id": self.check_id, "status": self.status, "evidence_digest": self.evidence_digest}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str) -> Self:
|
|
item = _object(value, path, ("check_id", "status", "evidence_digest"))
|
|
return cls(
|
|
check_id=_string(item["check_id"], f"{path}.check_id"),
|
|
status=_string(item["status"], f"{path}.status"),
|
|
evidence_digest=_string(item["evidence_digest"], f"{path}.evidence_digest"),
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class OutputQuality:
|
|
status: str
|
|
checks: tuple[OutputQualityCheck, ...]
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "status", _enum(self.status, "$.output_quality.status", {"failed", "passed"}))
|
|
if type(self.checks) not in {tuple, list} or not self.checks:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.output_quality.checks", "must be a non-empty sequence")
|
|
checks: list[OutputQualityCheck] = []
|
|
for index, check in enumerate(self.checks):
|
|
if not isinstance(check, OutputQualityCheck):
|
|
_fail(ContractErrorCode.TYPE_ERROR, f"$.output_quality.checks[{index}]", "must be OutputQualityCheck")
|
|
checks.append(check)
|
|
if len({check.check_id for check in checks}) != len(checks):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.output_quality.checks", "duplicate check_id")
|
|
object.__setattr__(self, "checks", tuple(checks))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {"status": self.status, "checks": [check.to_dict() for check in self.checks]}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.output_quality") -> Self:
|
|
item = _object(value, path, ("status", "checks"))
|
|
raw_checks = _array(item["checks"], f"{path}.checks", minimum=1)
|
|
return cls(
|
|
status=_string(item["status"], f"{path}.status"),
|
|
checks=tuple(
|
|
OutputQualityCheck.from_dict(check, f"{path}.checks[{index}]")
|
|
for index, check in enumerate(raw_checks)
|
|
),
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class OutputCoverage:
|
|
status: str
|
|
expected_count: int
|
|
observed_count: int
|
|
unit: str
|
|
domain: str
|
|
evidence_digest: str
|
|
|
|
def __post_init__(self) -> None:
|
|
object.__setattr__(self, "status", _enum(self.status, "$.output_coverage.status", {"complete", "incomplete", "stale", "unavailable"}))
|
|
object.__setattr__(self, "expected_count", _safe_integer(self.expected_count, "$.output_coverage.expected_count", minimum=1))
|
|
object.__setattr__(self, "observed_count", _safe_integer(self.observed_count, "$.output_coverage.observed_count", minimum=0))
|
|
object.__setattr__(self, "unit", _string(self.unit, "$.output_coverage.unit", _FIELD_NAME))
|
|
object.__setattr__(self, "domain", _logical_id(self.domain, "$.output_coverage.domain"))
|
|
object.__setattr__(self, "evidence_digest", _digest(self.evidence_digest, "$.output_coverage.evidence_digest"))
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"status": self.status,
|
|
"expected_count": self.expected_count,
|
|
"observed_count": self.observed_count,
|
|
"unit": self.unit,
|
|
"domain": self.domain,
|
|
"evidence_digest": self.evidence_digest,
|
|
}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.output_coverage") -> Self:
|
|
item = _object(
|
|
value,
|
|
path,
|
|
("status", "expected_count", "observed_count", "unit", "domain", "evidence_digest"),
|
|
)
|
|
return cls(
|
|
status=_string(item["status"], f"{path}.status"),
|
|
expected_count=item["expected_count"],
|
|
observed_count=item["observed_count"],
|
|
unit=_string(item["unit"], f"{path}.unit"),
|
|
domain=_string(item["domain"], f"{path}.domain"),
|
|
evidence_digest=_string(item["evidence_digest"], f"{path}.evidence_digest"),
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True, init=False)
|
|
class OutputArtifactRef:
|
|
artifact_id: str
|
|
schema_digest: str
|
|
content_digest: str
|
|
|
|
@classmethod
|
|
def create(cls, *, schema_digest: str, content_digest: str) -> Self:
|
|
normalized_schema = _digest(schema_digest, "$.output_artifact_ref.schema_digest")
|
|
normalized_content = _digest(content_digest, "$.output_artifact_ref.content_digest")
|
|
payload = {"schema_digest": normalized_schema, "content_digest": normalized_content}
|
|
artifact_id = f"rhfactoroutputv1:sha256:{hashlib.sha256(canonical_json_bytes(payload)).hexdigest()}"
|
|
instance = object.__new__(cls)
|
|
object.__setattr__(instance, "artifact_id", artifact_id)
|
|
object.__setattr__(instance, "schema_digest", normalized_schema)
|
|
object.__setattr__(instance, "content_digest", normalized_content)
|
|
return instance
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"artifact_id": self.artifact_id,
|
|
"schema_digest": self.schema_digest,
|
|
"content_digest": self.content_digest,
|
|
}
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.output_artifact_ref") -> Self:
|
|
item = _object(value, path, ("artifact_id", "schema_digest", "content_digest"))
|
|
supplied = _string(item["artifact_id"], f"{path}.artifact_id", _OUTPUT_ARTIFACT_ID)
|
|
result = cls.create(
|
|
schema_digest=_string(item["schema_digest"], f"{path}.schema_digest"),
|
|
content_digest=_string(item["content_digest"], f"{path}.content_digest"),
|
|
)
|
|
if supplied != result.artifact_id:
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, f"{path}.artifact_id", "artifact digest binding mismatch")
|
|
return result
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class UpstreamEvidence:
|
|
snapshot_id: str
|
|
snapshot_content_digest: str
|
|
snapshot_manifest_digest: str
|
|
quality_status: str
|
|
quality_evidence_digests: tuple[str, ...]
|
|
qualification_status: str
|
|
qualification_policy_id: str
|
|
qualification_policy_version: str
|
|
qualification_evaluated_at: str
|
|
qualification_evidence_digest: str
|
|
foundation_id: str
|
|
foundation_evidence_scope: str
|
|
foundation_contract_evidence_digests: tuple[str, ...]
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"snapshot_id": self.snapshot_id,
|
|
"snapshot_content_digest": self.snapshot_content_digest,
|
|
"snapshot_manifest_digest": self.snapshot_manifest_digest,
|
|
"quality_status": self.quality_status,
|
|
"quality_evidence_digests": list(self.quality_evidence_digests),
|
|
"qualification_status": self.qualification_status,
|
|
"qualification_policy_id": self.qualification_policy_id,
|
|
"qualification_policy_version": self.qualification_policy_version,
|
|
"qualification_evaluated_at": self.qualification_evaluated_at,
|
|
"qualification_evidence_digest": self.qualification_evidence_digest,
|
|
"foundation_id": self.foundation_id,
|
|
"foundation_evidence_scope": self.foundation_evidence_scope,
|
|
"foundation_contract_evidence_digests": list(self.foundation_contract_evidence_digests),
|
|
}
|
|
|
|
@classmethod
|
|
def derive(cls, snapshot: DatasetSnapshotEnvelope, foundation: DataFoundationEnvelope) -> Self:
|
|
return cls(
|
|
snapshot_id=snapshot.snapshot_id,
|
|
snapshot_content_digest=snapshot.content_digest,
|
|
snapshot_manifest_digest=snapshot.manifest_digest,
|
|
quality_status=snapshot.quality_status,
|
|
quality_evidence_digests=snapshot.quality_evidence_digests,
|
|
qualification_status=snapshot.qualification_status,
|
|
qualification_policy_id=snapshot.qualification_policy_id,
|
|
qualification_policy_version=snapshot.qualification_policy_version,
|
|
qualification_evaluated_at=snapshot.qualification_evaluated_at,
|
|
qualification_evidence_digest=snapshot.qualification_evidence_digest,
|
|
foundation_id=foundation.foundation_id,
|
|
foundation_evidence_scope=foundation.evidence_scope,
|
|
foundation_contract_evidence_digests=foundation.contract_evidence_digests,
|
|
)
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, path: str = "$.upstream_evidence") -> Self:
|
|
item = _object(
|
|
value,
|
|
path,
|
|
(
|
|
"snapshot_id",
|
|
"snapshot_content_digest",
|
|
"snapshot_manifest_digest",
|
|
"quality_status",
|
|
"quality_evidence_digests",
|
|
"qualification_status",
|
|
"qualification_policy_id",
|
|
"qualification_policy_version",
|
|
"qualification_evaluated_at",
|
|
"qualification_evidence_digest",
|
|
"foundation_id",
|
|
"foundation_evidence_scope",
|
|
"foundation_contract_evidence_digests",
|
|
),
|
|
)
|
|
quality_evidence = _array(item["quality_evidence_digests"], f"{path}.quality_evidence_digests", minimum=1)
|
|
foundation_evidence = _array(item["foundation_contract_evidence_digests"], f"{path}.foundation_contract_evidence_digests", minimum=1)
|
|
return cls(
|
|
snapshot_id=_string(item["snapshot_id"], f"{path}.snapshot_id", _SNAPSHOT_ID),
|
|
snapshot_content_digest=_digest(item["snapshot_content_digest"], f"{path}.snapshot_content_digest"),
|
|
snapshot_manifest_digest=_digest(item["snapshot_manifest_digest"], f"{path}.snapshot_manifest_digest"),
|
|
quality_status=_enum(item["quality_status"], f"{path}.quality_status", {"failed", "passed"}),
|
|
quality_evidence_digests=tuple(_digest(digest, f"{path}.quality_evidence_digests[{index}]") for index, digest in enumerate(quality_evidence)),
|
|
qualification_status=_enum(item["qualification_status"], f"{path}.qualification_status", {"qualified", "rejected"}),
|
|
qualification_policy_id=_string(item["qualification_policy_id"], f"{path}.qualification_policy_id"),
|
|
qualification_policy_version=_semver(item["qualification_policy_version"], f"{path}.qualification_policy_version"),
|
|
qualification_evaluated_at=_string(item["qualification_evaluated_at"], f"{path}.qualification_evaluated_at", _UTC_INSTANT),
|
|
qualification_evidence_digest=_digest(item["qualification_evidence_digest"], f"{path}.qualification_evidence_digest"),
|
|
foundation_id=_string(item["foundation_id"], f"{path}.foundation_id", _FOUNDATION_ID),
|
|
foundation_evidence_scope=_enum(item["foundation_evidence_scope"], f"{path}.foundation_evidence_scope", {"real_data", "synthetic_fixture"}),
|
|
foundation_contract_evidence_digests=tuple(_digest(digest, f"{path}.foundation_contract_evidence_digests[{index}]") for index, digest in enumerate(foundation_evidence)),
|
|
)
|
|
|
|
|
|
def _canonical_evidence_bytes(value: Any, path: str) -> bytes:
|
|
if type(value) is not bytes:
|
|
_fail(ContractErrorCode.TYPE_ERROR, path, "must be canonical JSON bytes")
|
|
try:
|
|
loaded = json.loads(value, object_pairs_hook=_duplicate_key_pairs)
|
|
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
|
|
raise FactorContractError(ContractErrorCode.INVALID_FORMAT, path, "invalid canonical JSON bytes") from exc
|
|
_assert_canonical_profile(loaded, path)
|
|
if canonical_json_bytes(loaded) != value:
|
|
_fail(ContractErrorCode.INVALID_FORMAT, path, "bytes are not canonical JSON")
|
|
return value
|
|
|
|
|
|
def _normalize_bindings(bindings: Sequence[InputBinding]) -> tuple[InputBinding, ...]:
|
|
if type(bindings) not in {tuple, list}:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.input_bindings", "must be a sequence")
|
|
normalized: list[InputBinding] = []
|
|
for index, binding in enumerate(bindings):
|
|
if not isinstance(binding, InputBinding):
|
|
_fail(ContractErrorCode.TYPE_ERROR, f"$.input_bindings[{index}]", "must be InputBinding")
|
|
normalized.append(binding)
|
|
keys = [(binding.definition_id, binding.input_name) for binding in normalized]
|
|
if len(keys) != len(set(keys)):
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.input_bindings", "duplicate factor input binding")
|
|
return tuple(sorted(normalized, key=lambda binding: (binding.definition_id, binding.input_name)))
|
|
|
|
|
|
def _normalize_view_availability(values: Sequence[ViewAvailability]) -> tuple[ViewAvailability, ...]:
|
|
if type(values) not in {tuple, list}:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.view_availability", "must be a sequence")
|
|
normalized: list[ViewAvailability] = []
|
|
for index, value in enumerate(values):
|
|
if not isinstance(value, ViewAvailability):
|
|
_fail(ContractErrorCode.TYPE_ERROR, f"$.view_availability[{index}]", "must be ViewAvailability")
|
|
normalized.append(value)
|
|
identities = [value.view_ref_id for value in normalized]
|
|
if len(identities) != len(set(identities)):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.view_availability", "duplicate view_ref_id")
|
|
return tuple(sorted(normalized, key=lambda value: value.view_ref_id))
|
|
|
|
|
|
@dataclass(frozen=True, slots=True, init=False)
|
|
class FactorSetRef:
|
|
contract_name: str
|
|
schema_version: str
|
|
factor_set_id: str
|
|
definition_ids: tuple[str, ...]
|
|
dataset_snapshot_id: str
|
|
foundation_id: str
|
|
pit_cutoff: str
|
|
selected_view_ref_ids: tuple[str, ...]
|
|
input_bindings: tuple[InputBinding, ...]
|
|
upstream_evidence: UpstreamEvidence
|
|
view_availability: tuple[ViewAvailability, ...]
|
|
output_quality: OutputQuality
|
|
output_coverage: OutputCoverage
|
|
output_schema_digest: str
|
|
output_content_digest: str
|
|
output_artifact_ref: OutputArtifactRef
|
|
availability_mode: AvailabilityMode
|
|
historical_availability: HistoricalAvailability
|
|
evaluation_at: str
|
|
computed_at: str
|
|
artifact_available_at: str
|
|
producer: ProducerIdentity
|
|
code_revision: str
|
|
actor: ActorIdentity
|
|
correlation_id: str
|
|
causation: Causation
|
|
evidence_scope: str
|
|
decision_eligible: bool
|
|
payload_validation: PayloadValidation = field(compare=False)
|
|
_definitions: tuple[FactorDefinition, ...] = field(repr=False, compare=False)
|
|
_dataset_snapshot: DatasetSnapshotEnvelope = field(repr=False, compare=False)
|
|
_foundation: DataFoundationEnvelope = field(repr=False, compare=False)
|
|
|
|
@classmethod
|
|
def create(
|
|
cls,
|
|
*,
|
|
definitions: Sequence[FactorDefinition],
|
|
dataset_snapshot: DatasetSnapshotEnvelope,
|
|
foundation: DataFoundationEnvelope,
|
|
selected_view_ref_ids: Sequence[str],
|
|
input_bindings: Sequence[InputBinding],
|
|
view_availability: Sequence[ViewAvailability],
|
|
output_quality: OutputQuality,
|
|
output_coverage: OutputCoverage,
|
|
output_schema_bytes: bytes,
|
|
output_content_bytes: bytes,
|
|
output_artifact_ref: OutputArtifactRef,
|
|
availability_mode: AvailabilityMode | str,
|
|
evaluation_at: str,
|
|
computed_at: str,
|
|
artifact_available_at: str,
|
|
producer: ProducerIdentity,
|
|
code_revision: str,
|
|
actor: ActorIdentity,
|
|
correlation_id: str,
|
|
causation: Causation,
|
|
evidence_scope: str,
|
|
decision_eligible: bool,
|
|
parent: FactorSetRef | None = None,
|
|
) -> Self:
|
|
schema_bytes = _canonical_evidence_bytes(output_schema_bytes, "$.output_schema_bytes")
|
|
content_bytes = _canonical_evidence_bytes(output_content_bytes, "$.output_content_bytes")
|
|
return cls._build(
|
|
definitions=definitions,
|
|
dataset_snapshot=dataset_snapshot,
|
|
foundation=foundation,
|
|
selected_view_ref_ids=selected_view_ref_ids,
|
|
input_bindings=input_bindings,
|
|
upstream_evidence=None,
|
|
view_availability=view_availability,
|
|
output_quality=output_quality,
|
|
output_coverage=output_coverage,
|
|
output_schema_digest=_digest_bytes(schema_bytes),
|
|
output_content_digest=_digest_bytes(content_bytes),
|
|
output_artifact_ref=output_artifact_ref,
|
|
availability_mode=availability_mode,
|
|
historical_availability=None,
|
|
evaluation_at=evaluation_at,
|
|
computed_at=computed_at,
|
|
artifact_available_at=artifact_available_at,
|
|
producer=producer,
|
|
code_revision=code_revision,
|
|
actor=actor,
|
|
correlation_id=correlation_id,
|
|
causation=causation,
|
|
evidence_scope=evidence_scope,
|
|
decision_eligible=decision_eligible,
|
|
parent=parent,
|
|
supplied_factor_set_id=None,
|
|
payload_validation=PayloadValidation.PAYLOAD_REVALIDATED,
|
|
)
|
|
|
|
@classmethod
|
|
def _build(
|
|
cls,
|
|
*,
|
|
definitions: Sequence[FactorDefinition],
|
|
dataset_snapshot: Any,
|
|
foundation: Any,
|
|
selected_view_ref_ids: Sequence[str],
|
|
input_bindings: Sequence[InputBinding],
|
|
upstream_evidence: UpstreamEvidence | None,
|
|
view_availability: Sequence[ViewAvailability],
|
|
output_quality: Any,
|
|
output_coverage: Any,
|
|
output_schema_digest: Any,
|
|
output_content_digest: Any,
|
|
output_artifact_ref: Any,
|
|
availability_mode: Any,
|
|
historical_availability: Any,
|
|
evaluation_at: Any,
|
|
computed_at: Any,
|
|
artifact_available_at: Any,
|
|
producer: Any,
|
|
code_revision: Any,
|
|
actor: Any,
|
|
correlation_id: Any,
|
|
causation: Any,
|
|
evidence_scope: Any,
|
|
decision_eligible: Any,
|
|
parent: FactorSetRef | None,
|
|
supplied_factor_set_id: Any,
|
|
payload_validation: PayloadValidation,
|
|
) -> Self:
|
|
if not isinstance(dataset_snapshot, DatasetSnapshotEnvelope):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.dataset_snapshot", "complete validated DatasetSnapshotEnvelope required")
|
|
if not isinstance(foundation, DataFoundationEnvelope):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.foundation", "complete validated DataFoundationEnvelope required")
|
|
dataset_snapshot.require_qualified()
|
|
if foundation.dataset_snapshot_id != dataset_snapshot.snapshot_id:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.foundation.dataset_snapshot_id", "snapshot mismatch")
|
|
snapshot_pit = _parse_utc(dataset_snapshot.pit_cutoff, "$.dataset_snapshot.pit_cutoff")
|
|
foundation_pit = _parse_utc(foundation.pit_cutoff, "$.foundation.pit_cutoff")
|
|
if snapshot_pit > foundation_pit:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.pit_cutoff", "snapshot PIT exceeds Foundation PIT")
|
|
normalized_definitions = validate_factor_catalog(definitions)
|
|
definition_by_id = {definition.definition_id: definition for definition in normalized_definitions}
|
|
if type(selected_view_ref_ids) not in {tuple, list}:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.selected_view_ref_ids", "must be a sequence")
|
|
selected_views = tuple(
|
|
sorted(
|
|
_string(view_id, f"$.selected_view_ref_ids[{index}]", _VIEW_REF_ID)
|
|
for index, view_id in enumerate(selected_view_ref_ids)
|
|
)
|
|
)
|
|
if not selected_views or len(selected_views) != len(set(selected_views)):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.selected_view_ref_ids", "must be non-empty and unique")
|
|
if not set(selected_views).issubset(foundation.views):
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.selected_view_ref_ids", "view is not a Foundation member")
|
|
bindings = _normalize_bindings(input_bindings)
|
|
expected_binding_keys = {
|
|
(definition.definition_id, factor_input.input_name)
|
|
for definition in normalized_definitions
|
|
for factor_input in definition.inputs
|
|
}
|
|
actual_binding_keys = {(binding.definition_id, binding.input_name) for binding in bindings}
|
|
if actual_binding_keys != expected_binding_keys:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.input_bindings", "factor input closure mismatch")
|
|
consumed_views: set[str] = set()
|
|
for binding in bindings:
|
|
definition = definition_by_id[binding.definition_id]
|
|
declared_input = next(item for item in definition.inputs if item.input_name == binding.input_name)
|
|
if binding.view_ref_id not in selected_views:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.input_bindings[].view_ref_id", "binding uses an unselected view")
|
|
view = foundation.views[binding.view_ref_id]
|
|
if binding.schema_digest != declared_input.schema_digest or binding.schema_digest != view.schema_digest:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.input_bindings[].schema_digest", "definition/view schema mismatch")
|
|
consumed_views.add(binding.view_ref_id)
|
|
if consumed_views != set(selected_views):
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.selected_view_ref_ids", "unused selected view")
|
|
availability = _normalize_view_availability(view_availability)
|
|
if {item.view_ref_id for item in availability} != set(selected_views):
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.view_availability", "must cover selected views exactly")
|
|
expected_upstream = UpstreamEvidence.derive(dataset_snapshot, foundation)
|
|
if upstream_evidence is not None and upstream_evidence != expected_upstream:
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, "$.upstream_evidence", "does not match admitted upstream envelopes")
|
|
if not isinstance(output_quality, OutputQuality):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.output_quality", "must be OutputQuality")
|
|
if output_quality.status != "passed" or any(check.status != "passed" for check in output_quality.checks):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.output_quality", "successful FactorSetRef requires every check passed")
|
|
if not isinstance(output_coverage, OutputCoverage):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.output_coverage", "must be OutputCoverage")
|
|
if output_coverage.status != "complete" or output_coverage.observed_count != output_coverage.expected_count:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.output_coverage", "successful FactorSetRef requires complete coverage")
|
|
schema_digest = _digest(output_schema_digest, "$.output_schema_digest")
|
|
content_digest = _digest(output_content_digest, "$.output_content_digest")
|
|
if not isinstance(output_artifact_ref, OutputArtifactRef):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.output_artifact_ref", "immutable artifact reference is required")
|
|
if output_artifact_ref.schema_digest != schema_digest or output_artifact_ref.content_digest != content_digest:
|
|
_fail(ContractErrorCode.ARTIFACT_MISMATCH, "$.output_artifact_ref", "artifact digests do not match output bytes")
|
|
try:
|
|
mode = availability_mode if isinstance(availability_mode, AvailabilityMode) else AvailabilityMode(availability_mode)
|
|
except (TypeError, ValueError) as exc:
|
|
raise FactorContractError(ContractErrorCode.INVALID_VALUE, "$.availability_mode", "unsupported availability mode") from exc
|
|
expected_historical = (
|
|
HistoricalAvailability.DECLARED_AS_AVAILABLE
|
|
if mode is AvailabilityMode.AS_AVAILABLE
|
|
else HistoricalAvailability.NOT_ESTABLISHED
|
|
)
|
|
if historical_availability is not None:
|
|
try:
|
|
supplied_historical = (
|
|
historical_availability
|
|
if isinstance(historical_availability, HistoricalAvailability)
|
|
else HistoricalAvailability(historical_availability)
|
|
)
|
|
except (TypeError, ValueError) as exc:
|
|
raise FactorContractError(ContractErrorCode.INVALID_VALUE, "$.historical_availability", "unsupported historical claim") from exc
|
|
if supplied_historical is not expected_historical:
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, "$.historical_availability", "claim is derived from availability_mode")
|
|
evaluation_text = _normalize_new_instant(evaluation_at, "$.evaluation_at")
|
|
computed_text = _normalize_new_instant(computed_at, "$.computed_at")
|
|
artifact_text = _normalize_new_instant(artifact_available_at, "$.artifact_available_at")
|
|
evaluation_time = _parse_utc(evaluation_text, "$.evaluation_at")
|
|
computed_time = _parse_utc(computed_text, "$.computed_at")
|
|
artifact_time = _parse_utc(artifact_text, "$.artifact_available_at")
|
|
if not dataset_snapshot.knowledge_start <= dataset_snapshot.knowledge_end <= snapshot_pit <= foundation_pit <= evaluation_time:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.evaluation_at", "knowledge <= snapshot PIT <= Foundation PIT <= evaluation required")
|
|
for definition in normalized_definitions:
|
|
if not _parse_utc(definition.valid_from, "$.definitions[].valid_from") <= evaluation_time < _parse_utc(
|
|
definition.valid_until, "$.definitions[].valid_until"
|
|
):
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.definitions[].validity", "definition is not valid at evaluation")
|
|
view_times = [_parse_utc(item.available_at, "$.view_availability[].available_at") for item in availability]
|
|
if mode is AvailabilityMode.AS_AVAILABLE:
|
|
if dataset_snapshot.published_at > foundation_pit:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.dataset_snapshot.descriptor.published_at", "publication is after Foundation PIT")
|
|
if any(view_time > foundation_pit for view_time in view_times):
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.view_availability", "view availability is after Foundation PIT")
|
|
if not all(source_time <= computed_time for source_time in [dataset_snapshot.published_at, *view_times]):
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.computed_at", "computation precedes source/view availability")
|
|
if not computed_time <= artifact_time <= evaluation_time:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.artifact_available_at", "computed <= artifact <= evaluation required")
|
|
else:
|
|
if dataset_snapshot.published_at > computed_time or any(view_time > computed_time for view_time in view_times):
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.computed_at", "replay computation precedes actual source/view availability")
|
|
if not evaluation_time <= computed_time <= artifact_time:
|
|
_fail(ContractErrorCode.TIME_ORDER_VIOLATION, "$.computed_at", "evaluation <= computed <= artifact required for replay")
|
|
if not isinstance(producer, ProducerIdentity):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.producer", "must be ProducerIdentity")
|
|
if producer.id != "quant_engine":
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.producer.id", "FactorSet producer must be quant_engine")
|
|
normalized_revision = _git_revision(code_revision, "$.code_revision")
|
|
if not isinstance(actor, ActorIdentity):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.actor", "must be ActorIdentity")
|
|
normalized_correlation = _logical_id(correlation_id, "$.correlation_id")
|
|
if not isinstance(causation, Causation):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.causation", "must be Causation")
|
|
if causation.kind == "foundation":
|
|
if causation.id != foundation.foundation_id:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.causation.id", "must bind exact input Foundation")
|
|
if parent is not None:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.causation", "Foundation cause cannot carry a FactorSet parent")
|
|
else:
|
|
if parent is None:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.causation", "FactorSet cause requires the exact parent object")
|
|
if causation.id != parent.factor_set_id:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.causation.id", "parent factor_set_id mismatch")
|
|
if normalized_correlation != parent.correlation_id:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.correlation_id", "parent correlation mismatch")
|
|
normalized_scope = _enum(evidence_scope, "$.evidence_scope", {"real_data", "synthetic_fixture"})
|
|
if normalized_scope == "real_data" and (
|
|
foundation.evidence_scope != "real_data" or foundation.real_data_status != "validated"
|
|
):
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, "$.evidence_scope", "real-data evidence is not established upstream")
|
|
if type(decision_eligible) is not bool:
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.decision_eligible", "must be boolean")
|
|
if decision_eligible:
|
|
_fail(ContractErrorCode.READINESS_ESCALATION, "$.decision_eligible", "S2 computation artifacts are never decision eligible")
|
|
payload = {
|
|
"contract_name": "researchhub.factor-set-ref",
|
|
"schema_version": FACTOR_CONTRACT_VERSION,
|
|
"definition_ids": [definition.definition_id for definition in normalized_definitions],
|
|
"dataset_snapshot_id": dataset_snapshot.snapshot_id,
|
|
"foundation_id": foundation.foundation_id,
|
|
"pit_cutoff": foundation.pit_cutoff,
|
|
"selected_view_ref_ids": list(selected_views),
|
|
"input_bindings": [binding.to_dict() for binding in bindings],
|
|
"upstream_evidence": expected_upstream.to_dict(),
|
|
"view_availability": [item.to_dict() for item in availability],
|
|
"output_quality": output_quality.to_dict(),
|
|
"output_coverage": output_coverage.to_dict(),
|
|
"output_schema_digest": schema_digest,
|
|
"output_content_digest": content_digest,
|
|
"output_artifact_ref": output_artifact_ref.to_dict(),
|
|
"availability_mode": mode.value,
|
|
"historical_availability": expected_historical.value,
|
|
"evaluation_at": evaluation_text,
|
|
"computed_at": computed_text,
|
|
"artifact_available_at": artifact_text,
|
|
"producer": producer.to_dict(),
|
|
"code_revision": normalized_revision,
|
|
"actor": actor.to_dict(),
|
|
"correlation_id": normalized_correlation,
|
|
"causation": causation.to_dict(),
|
|
"evidence_scope": normalized_scope,
|
|
"decision_eligible": False,
|
|
}
|
|
expected_id = f"rhfactorsetv1:sha256:{hashlib.sha256(canonical_json_bytes(payload)).hexdigest()}"
|
|
if causation.kind == "factor_set" and causation.id == expected_id:
|
|
_fail(ContractErrorCode.LINEAGE_VIOLATION, "$.causation.id", "self parent is forbidden")
|
|
if supplied_factor_set_id is not None:
|
|
supplied = _string(supplied_factor_set_id, "$.factor_set_id", _FACTOR_SET_ID)
|
|
if supplied != expected_id:
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, "$.factor_set_id", "does not identify complete normalized FactorSetRef")
|
|
values: dict[str, Any] = {
|
|
**payload,
|
|
"factor_set_id": expected_id,
|
|
"definition_ids": tuple(
|
|
definition.definition_id for definition in normalized_definitions
|
|
),
|
|
"selected_view_ref_ids": selected_views,
|
|
"input_bindings": bindings,
|
|
"upstream_evidence": expected_upstream,
|
|
"view_availability": availability,
|
|
"output_quality": output_quality,
|
|
"output_coverage": output_coverage,
|
|
"output_artifact_ref": output_artifact_ref,
|
|
"availability_mode": mode,
|
|
"historical_availability": expected_historical,
|
|
"producer": producer,
|
|
"actor": actor,
|
|
"causation": causation,
|
|
"payload_validation": payload_validation,
|
|
"_definitions": normalized_definitions,
|
|
"_dataset_snapshot": dataset_snapshot,
|
|
"_foundation": foundation,
|
|
}
|
|
instance = object.__new__(cls)
|
|
for name, value in values.items():
|
|
object.__setattr__(instance, name, value)
|
|
return instance
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"contract_name": self.contract_name,
|
|
"schema_version": self.schema_version,
|
|
"factor_set_id": self.factor_set_id,
|
|
"definition_ids": list(self.definition_ids),
|
|
"dataset_snapshot_id": self.dataset_snapshot_id,
|
|
"foundation_id": self.foundation_id,
|
|
"pit_cutoff": self.pit_cutoff,
|
|
"selected_view_ref_ids": list(self.selected_view_ref_ids),
|
|
"input_bindings": [binding.to_dict() for binding in self.input_bindings],
|
|
"upstream_evidence": self.upstream_evidence.to_dict(),
|
|
"view_availability": [item.to_dict() for item in self.view_availability],
|
|
"output_quality": self.output_quality.to_dict(),
|
|
"output_coverage": self.output_coverage.to_dict(),
|
|
"output_schema_digest": self.output_schema_digest,
|
|
"output_content_digest": self.output_content_digest,
|
|
"output_artifact_ref": self.output_artifact_ref.to_dict(),
|
|
"availability_mode": self.availability_mode.value,
|
|
"historical_availability": self.historical_availability.value,
|
|
"evaluation_at": self.evaluation_at,
|
|
"computed_at": self.computed_at,
|
|
"artifact_available_at": self.artifact_available_at,
|
|
"producer": self.producer.to_dict(),
|
|
"code_revision": self.code_revision,
|
|
"actor": self.actor.to_dict(),
|
|
"correlation_id": self.correlation_id,
|
|
"causation": self.causation.to_dict(),
|
|
"evidence_scope": self.evidence_scope,
|
|
"decision_eligible": self.decision_eligible,
|
|
}
|
|
|
|
def to_json(self) -> str:
|
|
return canonical_json(self.to_dict())
|
|
|
|
@classmethod
|
|
def from_dict(
|
|
cls,
|
|
value: Any,
|
|
*,
|
|
definitions: Sequence[FactorDefinition],
|
|
dataset_snapshot: DatasetSnapshotEnvelope,
|
|
foundation: DataFoundationEnvelope,
|
|
parent: FactorSetRef | None = None,
|
|
output_schema_bytes: bytes | None = None,
|
|
output_content_bytes: bytes | None = None,
|
|
) -> Self:
|
|
item = _object(
|
|
value,
|
|
"$",
|
|
(
|
|
"contract_name",
|
|
"schema_version",
|
|
"factor_set_id",
|
|
"definition_ids",
|
|
"dataset_snapshot_id",
|
|
"foundation_id",
|
|
"pit_cutoff",
|
|
"selected_view_ref_ids",
|
|
"input_bindings",
|
|
"upstream_evidence",
|
|
"view_availability",
|
|
"output_quality",
|
|
"output_coverage",
|
|
"output_schema_digest",
|
|
"output_content_digest",
|
|
"output_artifact_ref",
|
|
"availability_mode",
|
|
"historical_availability",
|
|
"evaluation_at",
|
|
"computed_at",
|
|
"artifact_available_at",
|
|
"producer",
|
|
"code_revision",
|
|
"actor",
|
|
"correlation_id",
|
|
"causation",
|
|
"evidence_scope",
|
|
"decision_eligible",
|
|
),
|
|
)
|
|
if item["contract_name"] != "researchhub.factor-set-ref":
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.contract_name", "unsupported contract")
|
|
if item["schema_version"] != FACTOR_CONTRACT_VERSION:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.schema_version", "unsupported version")
|
|
if item["dataset_snapshot_id"] != dataset_snapshot.snapshot_id:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.dataset_snapshot_id", "external envelope mismatch")
|
|
if item["foundation_id"] != foundation.foundation_id or item["pit_cutoff"] != foundation.pit_cutoff:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.foundation_id", "external Foundation mismatch")
|
|
raw_definition_ids = _array(item["definition_ids"], "$.definition_ids", minimum=1)
|
|
definition_ids = [
|
|
_string(definition_id, f"$.definition_ids[{index}]")
|
|
for index, definition_id in enumerate(raw_definition_ids)
|
|
]
|
|
if len(definition_ids) != len(set(definition_ids)):
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.definition_ids", "items must be unique")
|
|
if set(definition_ids) != {definition.definition_id for definition in definitions}:
|
|
_fail(ContractErrorCode.INPUT_CLOSURE_VIOLATION, "$.definition_ids", "external definitions mismatch")
|
|
raw_views = _array(item["selected_view_ref_ids"], "$.selected_view_ref_ids", minimum=1, unique=True)
|
|
raw_bindings = _array(item["input_bindings"], "$.input_bindings", minimum=1, unique=True)
|
|
raw_availability = _array(item["view_availability"], "$.view_availability", minimum=1, unique=True)
|
|
if (output_schema_bytes is None) != (output_content_bytes is None):
|
|
_fail(ContractErrorCode.ARTIFACT_MISMATCH, "$.output_artifact_ref", "both payload byte sets are required together")
|
|
validation = PayloadValidation.REFERENCE_ONLY
|
|
output_schema_digest = item["output_schema_digest"]
|
|
output_content_digest = item["output_content_digest"]
|
|
if output_schema_bytes is not None and output_content_bytes is not None:
|
|
output_schema_digest = _digest_bytes(_canonical_evidence_bytes(output_schema_bytes, "$.output_schema_bytes"))
|
|
output_content_digest = _digest_bytes(_canonical_evidence_bytes(output_content_bytes, "$.output_content_bytes"))
|
|
if output_schema_digest != item["output_schema_digest"] or output_content_digest != item["output_content_digest"]:
|
|
_fail(ContractErrorCode.ARTIFACT_MISMATCH, "$.output_artifact_ref", "supplied payload bytes do not match references")
|
|
validation = PayloadValidation.PAYLOAD_REVALIDATED
|
|
return cls._build(
|
|
definitions=definitions,
|
|
dataset_snapshot=dataset_snapshot,
|
|
foundation=foundation,
|
|
selected_view_ref_ids=tuple(raw_views),
|
|
input_bindings=tuple(InputBinding.from_dict(binding, f"$.input_bindings[{index}]") for index, binding in enumerate(raw_bindings)),
|
|
upstream_evidence=UpstreamEvidence.from_dict(item["upstream_evidence"]),
|
|
view_availability=tuple(ViewAvailability.from_dict(availability, f"$.view_availability[{index}]") for index, availability in enumerate(raw_availability)),
|
|
output_quality=OutputQuality.from_dict(item["output_quality"]),
|
|
output_coverage=OutputCoverage.from_dict(item["output_coverage"]),
|
|
output_schema_digest=output_schema_digest,
|
|
output_content_digest=output_content_digest,
|
|
output_artifact_ref=OutputArtifactRef.from_dict(item["output_artifact_ref"]),
|
|
availability_mode=item["availability_mode"],
|
|
historical_availability=item["historical_availability"],
|
|
evaluation_at=item["evaluation_at"],
|
|
computed_at=item["computed_at"],
|
|
artifact_available_at=item["artifact_available_at"],
|
|
producer=ProducerIdentity.from_dict(item["producer"]),
|
|
code_revision=item["code_revision"],
|
|
actor=ActorIdentity.from_dict(item["actor"]),
|
|
correlation_id=item["correlation_id"],
|
|
causation=Causation.from_dict(item["causation"]),
|
|
evidence_scope=item["evidence_scope"],
|
|
decision_eligible=item["decision_eligible"],
|
|
parent=parent,
|
|
supplied_factor_set_id=item["factor_set_id"],
|
|
payload_validation=validation,
|
|
)
|
|
|
|
@classmethod
|
|
def from_json(
|
|
cls,
|
|
value: str | bytes,
|
|
**kwargs: Any,
|
|
) -> Self:
|
|
return cls.from_dict(_parse_json_object(value, "$"), **kwargs)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True, init=False)
|
|
class LegacyFactorBinding:
|
|
contract_name: str
|
|
schema_version: str
|
|
binding_id: str
|
|
definition_id: str
|
|
legacy_factor_id: str
|
|
legacy_version: str
|
|
legacy_definition_sha256: str
|
|
legacy_dataset_schema_version: str
|
|
canonical_input_schema_digest: str
|
|
correspondence_evidence_digest: str
|
|
|
|
@classmethod
|
|
def create(
|
|
cls,
|
|
*,
|
|
definition: FactorDefinition,
|
|
legacy_factor_id: str,
|
|
legacy_version: str,
|
|
legacy_definition_sha256: str,
|
|
legacy_dataset_schema_version: str,
|
|
canonical_input_schema_digest: str,
|
|
correspondence_evidence_digest: str,
|
|
) -> Self:
|
|
if not isinstance(definition, FactorDefinition):
|
|
_fail(ContractErrorCode.TYPE_ERROR, "$.definition", "complete FactorDefinition required")
|
|
normalized_canonical_digest = _digest(canonical_input_schema_digest, "$.canonical_input_schema_digest")
|
|
if normalized_canonical_digest != definition.input_schema_digest:
|
|
_fail(ContractErrorCode.LEGACY_BINDING_MISMATCH, "$.canonical_input_schema_digest", "definition input schema mismatch")
|
|
payload = {
|
|
"contract_name": "researchhub.legacy-factor-binding",
|
|
"schema_version": FACTOR_CONTRACT_VERSION,
|
|
"definition_id": definition.definition_id,
|
|
"legacy_factor_id": _logical_id(legacy_factor_id, "$.legacy_factor_id"),
|
|
"legacy_version": _semver(legacy_version, "$.legacy_version"),
|
|
"legacy_definition_sha256": _string(legacy_definition_sha256, "$.legacy_definition_sha256", _BARE_SHA256),
|
|
"legacy_dataset_schema_version": _semver(legacy_dataset_schema_version, "$.legacy_dataset_schema_version"),
|
|
"canonical_input_schema_digest": normalized_canonical_digest,
|
|
"correspondence_evidence_digest": _digest(correspondence_evidence_digest, "$.correspondence_evidence_digest"),
|
|
}
|
|
binding_id = f"rhlegacyfactorv1:sha256:{hashlib.sha256(canonical_json_bytes(payload)).hexdigest()}"
|
|
instance = object.__new__(cls)
|
|
for name, value in {**payload, "binding_id": binding_id}.items():
|
|
object.__setattr__(instance, name, value)
|
|
return instance
|
|
|
|
def to_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"contract_name": self.contract_name,
|
|
"schema_version": self.schema_version,
|
|
"binding_id": self.binding_id,
|
|
"definition_id": self.definition_id,
|
|
"legacy_factor_id": self.legacy_factor_id,
|
|
"legacy_version": self.legacy_version,
|
|
"legacy_definition_sha256": self.legacy_definition_sha256,
|
|
"legacy_dataset_schema_version": self.legacy_dataset_schema_version,
|
|
"canonical_input_schema_digest": self.canonical_input_schema_digest,
|
|
"correspondence_evidence_digest": self.correspondence_evidence_digest,
|
|
}
|
|
|
|
def to_json(self) -> str:
|
|
return canonical_json(self.to_dict())
|
|
|
|
@classmethod
|
|
def from_dict(cls, value: Any, *, definition: FactorDefinition) -> Self:
|
|
item = _object(
|
|
value,
|
|
"$",
|
|
(
|
|
"contract_name",
|
|
"schema_version",
|
|
"binding_id",
|
|
"definition_id",
|
|
"legacy_factor_id",
|
|
"legacy_version",
|
|
"legacy_definition_sha256",
|
|
"legacy_dataset_schema_version",
|
|
"canonical_input_schema_digest",
|
|
"correspondence_evidence_digest",
|
|
),
|
|
)
|
|
if item["contract_name"] != "researchhub.legacy-factor-binding" or item["schema_version"] != FACTOR_CONTRACT_VERSION:
|
|
_fail(ContractErrorCode.INVALID_VALUE, "$.contract_name", "unsupported legacy binding contract")
|
|
if item["definition_id"] != definition.definition_id:
|
|
_fail(ContractErrorCode.LEGACY_BINDING_MISMATCH, "$.definition_id", "definition mismatch")
|
|
result = cls.create(
|
|
definition=definition,
|
|
legacy_factor_id=item["legacy_factor_id"],
|
|
legacy_version=item["legacy_version"],
|
|
legacy_definition_sha256=item["legacy_definition_sha256"],
|
|
legacy_dataset_schema_version=item["legacy_dataset_schema_version"],
|
|
canonical_input_schema_digest=item["canonical_input_schema_digest"],
|
|
correspondence_evidence_digest=item["correspondence_evidence_digest"],
|
|
)
|
|
supplied = _string(item["binding_id"], "$.binding_id", _LEGACY_BINDING_ID)
|
|
if supplied != result.binding_id:
|
|
_fail(ContractErrorCode.IDENTITY_MISMATCH, "$.binding_id", "legacy binding digest mismatch")
|
|
return result
|
|
|
|
@classmethod
|
|
def from_json(cls, value: str | bytes, *, definition: FactorDefinition) -> Self:
|
|
return cls.from_dict(_parse_json_object(value, "$"), definition=definition)
|
|
|
|
|
|
__all__ = [
|
|
"ActorIdentity",
|
|
"AvailabilityMode",
|
|
"ContractErrorCode",
|
|
"Causation",
|
|
"DataFoundationEnvelope",
|
|
"DatasetSnapshotEnvelope",
|
|
"FACTOR_CONTRACT_VERSION",
|
|
"FactorContractError",
|
|
"FactorDefinition",
|
|
"FactorInput",
|
|
"FactorSetRef",
|
|
"HistoricalAvailability",
|
|
"InputBinding",
|
|
"LegacyFactorBinding",
|
|
"OutputArtifactRef",
|
|
"OutputCoverage",
|
|
"OutputQuality",
|
|
"OutputQualityCheck",
|
|
"PayloadValidation",
|
|
"ProducerIdentity",
|
|
"TypedParameter",
|
|
"UpstreamEvidence",
|
|
"ViewAvailability",
|
|
"canonical_json",
|
|
"canonical_json_bytes",
|
|
"factor_definition_from_alpha158",
|
|
"factor_input_schema_digest",
|
|
"validate_factor_catalog",
|
|
]
|