diff --git a/docs/rates-context-observer-research.md b/docs/rates-context-observer-research.md index bc15ce1..7451ab1 100644 --- a/docs/rates-context-observer-research.md +++ b/docs/rates-context-observer-research.md @@ -1,6 +1,6 @@ # Ten-year rates context observation research -Status: `UNVALIDATED_RESEARCH`. This is one source-agnostic, standard-library pure observation function. It is not a downloader, enabled plugin, strategy policy, or production adoption. Existing modules, package exports, runner/catalog entries, wide external-context CSVs, defaults, dependencies, and runtime pins are unchanged. +Status: `UNVALIDATED_RESEARCH`. This module has source-agnostic, standard-library pure observation entry points. It is not a downloader, enabled plugin, strategy policy, or production adoption. Existing modules, package exports, runner/catalog entries, wide external-context CSVs, defaults, dependencies, and runtime pins are unchanged. [中文合同](rates-context-observer-research.zh-CN.md) @@ -10,7 +10,7 @@ Status: `UNVALIDATED_RESEARCH`. This is one source-agnostic, standard-library pu The existing wide CSV merge coerces each column to numeric values. It cannot preserve this contract's source, revision, or publication metadata; do not insert these records into that merge and assume provenance survives. This module neither extends `DEFAULT_FRED_SERIES` nor modifies macro watch/actionable scoring. -## Input and configuration +## Existing v1 input and configuration The input has exactly `schema_version=qsl.rates-context-input.research.v1`, `decision_at`, and `series`. `decision_at` is the actual evaluation/receipt decision time declared by the caller, with an explicit ISO timezone. `series` contains only the optional roles `nominal_10y`, `real_10y`, and `breakeven_10y`; a missing role yields explicit unknown. All three roles are ten-year measurements, not ETF prices or another tenor relabeled as ten-year. @@ -72,3 +72,97 @@ Focused tests can run with `python -m unittest discover -s tests -p test_rates_c Synthetic coverage includes percent-to-bp arithmetic, independent breakeven, individual source delays, negative yields, future-row/revision invariance, missing publication/receipt evidence, today's import of old data, nonfinite/Boolean/string values, wrong units, mixed sources/methodologies, reversed dates, visible duplicate revisions, wrong/missing endpoints, stale observations, invalid configuration, strict JSON, input immutability, and absent trading fields. This is a preparation step. A future authorized collector/consumer integration must retain the row-level evidence rather than using the lossy wide CSV, then validate genuine source coverage/availability, strategy-specific economic usefulness, and approved consumption. Historical breadth membership/prices and NDX participation are separate gaps; this phase implements neither breadth nor an NDX proxy. No runtime adoption follows from a pure function or passing synthetic tests. + +## Explicit forward-known v2 entry point + +build_rates_context_observation_v2(snapshot, config) is a separate pure entry +point in the same module. It accepts only qsl.rates-context-input.research.v2 +and returns qsl.rates-context-observation.research.v2. The original v1 entry +point, constants, input/output semantics and rejection behavior remain unchanged. +Neither entry point automatically upgrades the other's input. V2 never substitutes +first-seen for v1 available_at. + +The v2 root has exactly schema_version, availability_basis, collector_id, +decision_at and series. availability_basis is collector_first_seen. collector_id +identifies the collector to which the declared first-seen time belongs: a nonempty, +whitespace-trimmed string of at most 256 characters. Existing source/series/basis, +percent-unit roles and explicit v1 window/age configuration are reused. There is +no MSS family, catalog or runtime consumer registration. + +Each v2 row has exactly: + +- observation_date, value, revision_id +- source_published_at: explicitly null, never omitted or inferred +- first_seen_at: when the named collector first completely received this declared semantic row version +- received_at: when the research consumer received that fixed version +- capture_sha256: the declared hash of the original complete response retained when the row version was first seen +- row_sha256: the producer's declared immutable row-content reference + +Both hashes require lowercase 64-hex syntax. The observer does not read the +response, recompute either hash, authenticate origin or clocks, establish an +earliest capture, or verify the producer's hashing algorithm. These are references +and consistency checks, not historical PIT proof. The collector retains original +bytes/capture records externally. Repeated downloads must not refresh an old +first_seen or replace its first-capture reference with the latest whole-file hash. +Corrections require separately retained versions; old accepted decisions are not +rewritten. This stateless function cannot enforce those rules across calls. + +### Time projection and visible conflicts + +V2 timestamps require YYYY-MM-DDTHH:MM:SS, optionally 1–6 fractional second +digits, followed by Z or ±HH:MM (offset hours <=23, minutes <=59). Greater +precision, comma fractions and second-bearing offsets are unsupported and remain +unknown; times are never truncated or rounded. Valid timestamps normalize to UTC. +Selected rows require +first_seen_at <= received_at <= decision_at; the UTC first-seen date cannot +precede the observation date. known_at is the later first-seen/receipt time, +which equals receipt under valid ordering. It is not a publication timestamp. + +Dates outside the explicit window, future observation dates, and rows whose +first-seen or receipt is later than the decision are excluded before non-selector +fields. Hidden rows cannot affect counts, validation, identities, endpoint records +or numeric results. An unparseable selector cannot prove invisibility and produces +unknown unless another valid selector has already excluded the row. + +Visible rows retain strict chronological order. Two rows for one date, including +exact repeats, are ambiguous: no sorting, deduplication, latest-wins or revision +selection. A duplicated endpoint has no selected endpoint metadata. Conflicting +content or timing under one observation/revision identity is unknown. Reuse of one +row hash for different observation dates or values within a series is unknown. +One capture hash may legitimately cover multiple different rows. + +An old observation first captured today may be known today when its explicit +window/age policy permits it. It is never visible to an earlier decision. Age stays +observation-date age, not capture/receipt age; today's import does not refresh it. + +### Output and assurance + +V2 preserves collector/basis and each unambiguous endpoint's null publication, +first-seen, receipt, known-at, revision and hash references. There is no available_at +field. first_seen_delay_calendar_days compares date labels; +consumer_receipt_lag_seconds measures first-seen-to-consumer receipt. Neither is +source-publication latency or market-close-to-publication latency. + +Missing/invalid fields, identities, times, nonfinite values, visible conflicts, +wrong roles/bases/units, stale observations or unavailable exact endpoints produce +unknown and no numeric change. Invalid configuration retains ContractError. +AI, opportunity and control fields are outside this exact input contract. + +The same-source nominal/real approximate spread remains separate from independent +reported breakeven. Missing breakeven does not erase valid nominal/real declaration +facts or a valid pair difference; overall quality remains unknown. Every output +retains UNVALIDATED_RESEARCH, CALLER_DECLARATIONS_AND_CONSISTENCY_ONLY, +historical_pit_verified=false, backtest_eligible=false and +position_control_allowed=false. + +This is offline contract preparation. Actual capture, source rights, retained +source bytes, genuine first-seen clocks, forward history, economic qualification +and approved consumption remain unverified. A later separately reviewed MSS +adapter can consume fixed v2 fields without changing v1. Real collection/adoption +need their own evidence; pre-collection historical availability is not recreated. + +The existing focused unittest command runs original v1 plus synthetic v2 cases: +version/collector gates, arrival boundaries, old-date visibility without backdating, +stale age, hidden future/revision invariance, visible ambiguity, hash-reference +consistency, partial series, absent independent breakeven, strict JSON and unchanged +inputs. Pure tests do not prove external first-seen persistence. diff --git a/docs/rates-context-observer-research.zh-CN.md b/docs/rates-context-observer-research.zh-CN.md index ad9e44c..08c1a7c 100644 --- a/docs/rates-context-observer-research.zh-CN.md +++ b/docs/rates-context-observer-research.zh-CN.md @@ -1,6 +1,6 @@ # 十年利率背景观察研究 -状态:`UNVALIDATED_RESEARCH`。本期只有一个来源无关、标准库纯观察函数。不是下载器、已启用插件、策略政策或生产采用。既有模块、包级导出、runner/catalog、外部背景宽表CSV、默认配置、依赖和runtime pin保持原样。 +状态:`UNVALIDATED_RESEARCH`。保留原来源无关、标准库纯 v1 观察入口,并新增独立显式 v2 入口。不是下载器、已启用插件、策略政策或生产采用。既有模块、包级导出、runner/catalog、外部背景宽表CSV、默认配置、依赖和runtime pin保持原样。 [Full English contract](rates-context-observer-research.md) @@ -10,7 +10,7 @@ 原宽表CSV合并会将各列转成数值,不能保存本合同的来源、修订和发布时间。不能把新记录塞进该合并后声称来源证据仍完整。本模块不扩充 `DEFAULT_FRED_SERIES`,不修改macro的watch/actionable评分。 -## 输入与配置 +## 原 v1 输入与配置 输入恰有 `schema_version=qsl.rates-context-input.research.v1`、`decision_at`、`series`。`decision_at` 是调用方声明的实际评估/接收决策时点,须有明确ISO时区。`series` 仅允许可缺省的 `nominal_10y`、`real_10y`、`breakeven_10y` 三个角色;缺少角色明确unknown。三者必须是真正十年期测量,不能把ETF价格或其他期限重标为十年。 @@ -72,3 +72,79 @@ focused测试可用 `python -m unittest discover -s tests -p test_rates_context_ 合成覆盖percent→bp、独立breakeven、各源延迟、负利率、未来行/修订不变、缺公布/接收证据、今天导入旧数据、NaN/无穷/布尔/字符串、错误单位、混源/混口径、倒序/可见重复修订、错误/缺失端点、过时、非法配置、严格JSON、输入不变及无交易字段。 这只是准备步骤。未来获准采集/consumer接线须保留逐行证据,不能经过丢失metadata的宽表;另验真实来源覆盖/可得性、策略经济价值及获批消费。breadth历史membership/prices和NDX参与结构是独立缺口,本期没有实现breadth或NDX代理。纯函数或合成测试通过均不证明runtime采用。 + +## 显式 forward-known v2 入口 + +同模块新增独立纯入口 build_rates_context_observation_v2(snapshot, config), +仅接收 qsl.rates-context-input.research.v2,返回 +qsl.rates-context-observation.research.v2。原 v1 入口、常量、输入/输出语义 +及拒绝行为保持不变。两个入口不自动升级对方输入;v2 不把 first_seen +填入 v1 的 available_at。 + +v2 根字段恰为 schema_version、availability_basis、collector_id、decision_at、 +series。availability_basis 固定 collector_first_seen。collector_id 是 first_seen +所归属采集器的声明身份,须为无首尾空白、非空且不超过 256 字符的字符串。 +source/series/basis/percent 单位、三个角色和显式 v1 窗口/年龄配置继续沿用。 +本批不注册 MSS family、catalog 或运行消费者。 + +每个 v2 row 恰有: + +- observation_date、value、revision_id +- source_published_at:必须显式 null,不能省略或推算 +- first_seen_at:指定采集器首次完整收到该声明语义行版本的时间 +- received_at:研究 consumer 收到该固定版本的时间 +- capture_sha256:该行版本首次被看到时保留的原始完整响应之声明 hash +- row_sha256:producer 声明的不可变行内容引用 + +两个 hash 均要求小写 64 位十六进制形式。observer 不读取响应、不重算 hash、 +不认证来源或时钟、不确定“首次”真实性,也不验证 producer 的 hash 算法。 +这里是引用和一致性检查,不是历史 PIT 证据。collector 在外部保留原始 bytes +和捕获记录。重复下载不能刷新旧 first_seen 或把首次 capture 引用换成新版 +整文件 hash;更正须另留版本,不改写旧接受决策。无状态纯函数不能跨调用 +强制这些持久记录规则。 + +### 时间投影与可见冲突 + +v2 时间精确支持 YYYY-MM-DDTHH:MM:SS,可带 1–6 位小数秒,后接 Z 或 ±HH:MM。 +偏移小时不得大于 23、分钟不得大于 59;不支持更高精度、逗号小数或带秒的偏移。 +不截断/舍入时间;不支持的形式保持 unknown。合法时间归一 UTC。选中行满足 first_seen_at <= received_at <= decision_at; +first_seen 的 UTC 日期不得早于观察日。known_at 为 first_seen/received 较晚值, +合法顺序下等于 received,不是源发布时间。 + +先排除显式窗口外日期、未来观察日,以及 first_seen 或 received 晚于决策的行, +再处理非选择器字段。隐藏行不影响计数、校验、身份、端点记录和数值。 +无法解析的选择器不能证明不可见;若其他合法选择器尚未排除该行,则 unknown。 + +可见行保持严格日期顺序;同日两行,包括完全重复,也属于歧义。 +不排序、不去重、不 latest-wins、不自行选择修订。重复端点没有选中端点 metadata。 +同一观察/revision 身份下内容或时间冲突为 unknown;同一 series 中, +同一 row hash 对应不同观察日或数值也为 unknown。一个 capture hash 可以 +合法对应同份响应的多条不同记录。 + +旧观察今天首次采集后,可以在显式窗口/年龄允许时成为今天已知的信息, +不能成为更早决策可见的信息。年龄始终取观察日期,不按采集/接收时间刷新。 + +### 输出与保证范围 + +v2 保留 collector/basis 及每个无歧义端点的 null publication、first_seen、 +received、known_at、revision 和 hash 引用,不输出 available_at。 +first_seen_delay_calendar_days 仅比较日期标签;consumer_receipt_lag_seconds +表示首次采集到 consumer 接收的间隔。两者都不是源公布延迟或市场收盘至发布时间。 + +缺失/非法字段、身份、时间、非有限值、可见冲突、错误角色/basis/单位、 +过期或准确端点不可用,均 unknown 且不给数值变化。非法配置仍为 ContractError。 +AI、机会和控制字段不属于该精确合同。 + +同源 nominal/real 的 approximate spread 与独立 reported breakeven 分开。 +缺 breakeven 不抹掉两条有效声明事实或合法配对差值,整体 quality 仍 unknown。 +所有输出保持 UNVALIDATED_RESEARCH、CALLER_DECLARATIONS_AND_CONSISTENCY_ONLY、 +historical_pit_verified=false、backtest_eligible=false、position_control_allowed=false。 + +本批仅离线合同准备,未证明真实采集、来源权利、原始 bytes、可信 first_seen +时钟、前向历史、经济资格或获批消费。冻结 v2 后另行复审 MSS adapter 才接入, +不改旧 v1。真实采集/采用各需证据,不能还原采集前的历史可得性。 + +既有 focused unittest 命令同时运行原 v1 和合成 v2:版本/collector gate、 +到达边界、旧日期非回填可见性、旧数据年龄、未来/迟到修订不变性、 +可见版本歧义、hash 引用一致性、局部 series、缺独立 breakeven、严格 JSON +和输入不变。纯测试不能证明外部 first_seen 持久性。 diff --git a/src/quant_strategy_plugins/rates_context_observer_research.py b/src/quant_strategy_plugins/rates_context_observer_research.py index c6fdc43..eae8d34 100644 --- a/src/quant_strategy_plugins/rates_context_observer_research.py +++ b/src/quant_strategy_plugins/rates_context_observer_research.py @@ -233,3 +233,224 @@ def build_rates_context_observation(snapshot: Mapping, config: Mapping) -> dict: "max_observation_age_days": max_age, "series": observations, "approximate_yield_spread": spread, "quality": {"status": "declared_available" if all(item["status"] == "declared_available" for item in observations.values()) and spread["status"] == "declared_available" else "unknown"}} + + +# V2 is explicit and separate: v1's entry point, constants and semantics above +# remain unchanged. These fields describe collector knowledge, not publication. +FORWARD_INPUT_VERSION = "qsl.rates-context-input.research.v2" +FORWARD_OBSERVATION_VERSION = "qsl.rates-context-observation.research.v2" +_FORWARD_BASIS = "collector_first_seen" +_FORWARD_ROOT_KEYS = {"schema_version", "availability_basis", "collector_id", "decision_at", "series"} +_FORWARD_ROW_KEYS = {"observation_date", "value", "revision_id", "source_published_at", + "first_seen_at", "received_at", "capture_sha256", "row_sha256"} + + +def _forward_time(value): + # datetime.fromisoformat truncates excessive precision and normalizes some + # invalid offsets. V2 accepts an exact microsecond-representable grammar. + if not isinstance(value, str): + return None + if value.endswith("Z"): + body = value[:-1] + elif len(value) >= 6 and value[-6] in "+-" and value[-3] == ":": + offset = value[-5:-3] + value[-2:] + if not offset.isascii() or not offset.isdigit() or int(offset[:2]) > 23 or int(offset[2:]) > 59: + return None + body = value[:-6] + else: + return None + whole, separator, fraction = body.partition(".") + if len(whole) != 19 or (whole[4], whole[7], whole[10], whole[13], whole[16]) != ("-", "-", "T", ":", ":"): + return None + digits = whole[:4] + whole[5:7] + whole[8:10] + whole[11:13] + whole[14:16] + whole[17:19] + if not digits.isascii() or not digits.isdigit(): + return None + if separator and (not 1 <= len(fraction) <= 6 or not fraction.isascii() or not fraction.isdigit()): + return None + return _time(value) + + +def _forward_digest(value): + return value if isinstance(value, str) and len(value) == 64 and all(c in "0123456789abcdef" for c in value) else None + + +def _forward_endpoint(row): + if row is None: + return None + first_seen, received = _forward_time(row.get("first_seen_at")), _forward_time(row.get("received_at")) + known = received if first_seen is not None and received is not None and received >= first_seen else None + return {"observation_date": row["observation_date"], "source_published_at": None, + "first_seen_at": _stamp(first_seen), "received_at": _stamp(received), "known_at": _stamp(known), + "revision_id": _text(row.get("revision_id")), "capture_sha256": _forward_digest(row.get("capture_sha256")), + "row_sha256": _forward_digest(row.get("row_sha256"))} + + +def _observe_forward_series(series, role, start, end, decision, max_age, top_valid): + series = series if isinstance(series, Mapping) else {} + result = {"status": "unknown", "reason_codes": [], "source_id": _text(series.get("source_id")), + "series_id": _text(series.get("series_id")), "basis": _text(series.get("basis")), + "declared_unit": _text(series.get("unit")), "unit": "percent", "start_percent": None, + "end_percent": None, "change_bp": None, "start_observation": None, "end_observation": None, + "latest_observation_date": None, "observation_age_days": None, + "first_seen_delay_calendar_days": None, "consumer_receipt_lag_seconds": None, + "visible_observations": 0} + reasons = result["reason_codes"] + if not top_valid: + reasons.append("INPUT_INVALID") + return result + if not series: + reasons.append("SERIES_MISSING") + return result + if set(series) != _SERIES_KEYS or result["source_id"] is None or result["series_id"] is None: + reasons.append("SOURCE_IDENTITY_INVALID") + allowed_basis = {"reported_breakeven"} if role == "breakeven_10y" else _YIELD_BASES + if result["basis"] not in allowed_basis: + reasons.append("BASIS_INVALID") + if series.get("unit") != "percent": + reasons.append("UNIT_NOT_PERCENT") + rows = series.get("rows") + if not isinstance(rows, Sequence) or isinstance(rows, (str, bytes, bytearray)): + reasons.append("ROWS_INVALID") + return result + visible, revisions, content_ids = [], {}, {} + previous = None + for row in rows: + if not isinstance(row, Mapping) or (observed := _date(row.get("observation_date"))) is None: + reasons.append("ROW_DATE_INVALID") + continue + if observed < start or observed > end or observed > decision.date(): + continue + first_seen = _forward_time(row.get("first_seen_at")) + if first_seen is not None and first_seen > decision: + continue + received = _forward_time(row.get("received_at")) + if received is not None and received > decision: + continue + # Only selected rows can affect validation, counts, identity or output. + if previous is not None and observed <= previous: + reasons.append("ROW_ORDER_OR_REVISION_AMBIGUOUS") + previous = observed + visible.append(row) + if set(row) != _FORWARD_ROW_KEYS: + reasons.append("ROW_FIELDS_INVALID") + if "source_published_at" not in row or row.get("source_published_at") is not None: + reasons.append("SOURCE_PUBLICATION_MUST_BE_UNKNOWN") + if first_seen is None: + reasons.append("FIRST_SEEN_UNKNOWN") + elif first_seen.date() < observed: + reasons.append("FIRST_SEEN_BEFORE_OBSERVATION") + if received is None: + reasons.append("RECEIPT_UNKNOWN") + elif first_seen is not None and received < first_seen: + reasons.append("RECEIPT_BEFORE_FIRST_SEEN") + revision = _text(row.get("revision_id")) + if revision is None or revision == "latest": + reasons.append("REVISION_IDENTITY_INVALID") + capture_hash, row_hash = _forward_digest(row.get("capture_sha256")), _forward_digest(row.get("row_sha256")) + if capture_hash is None: + reasons.append("CAPTURE_IDENTITY_INVALID") + if row_hash is None: + reasons.append("ROW_CONTENT_IDENTITY_INVALID") + value = _number(row.get("value")) + if value is None: + reasons.append("VALUE_INVALID") + if revision is not None: + identity = (observed.isoformat(), revision) + content = (value, _stamp(first_seen), _stamp(received), capture_hash, row_hash) + if identity in revisions and revisions[identity] != content: + reasons.append("ROW_IDENTITY_CONFLICT") + revisions[identity] = content + if row_hash is not None: + content = (observed.isoformat(), value) + if row_hash in content_ids and content_ids[row_hash] != content: + reasons.append("ROW_IDENTITY_CONFLICT") + content_ids[row_hash] = content + result["visible_observations"] = len(visible) + first_rows = [row for row in visible if row["observation_date"] == start.isoformat()] + last_rows = [row for row in visible if row["observation_date"] == end.isoformat()] + # Ambiguity never selects first/last writer even for provenance display. + first = first_rows[0] if len(first_rows) == 1 else None + last = last_rows[0] if len(last_rows) == 1 else None + result["start_observation"], result["end_observation"] = _forward_endpoint(first), _forward_endpoint(last) + if not first_rows or not last_rows: + reasons.append("WINDOW_ENDPOINT_UNAVAILABLE") + if visible: + latest = max(_date(row["observation_date"]) for row in visible) + result["latest_observation_date"] = latest.isoformat() + result["observation_age_days"] = (decision.date() - latest).days + if result["observation_age_days"] > max_age: + reasons.append("OBSERVATION_STALE") + if last is not None: + first_seen, received = _forward_time(last.get("first_seen_at")), _forward_time(last.get("received_at")) + if first_seen is not None: + result["first_seen_delay_calendar_days"] = (first_seen.date() - end).days + if first_seen is not None and received is not None: + result["consumer_receipt_lag_seconds"] = (received - first_seen).total_seconds() + if not reasons: + first_value, last_value = _number(first["value"]), _number(last["value"]) + change = _rounded(100 * (last_value - first_value)) + if change is None: + reasons.append("CHANGE_NONFINITE") + else: + result.update(status="declared_available", start_percent=first_value, + end_percent=last_value, change_bp=change) + result["reason_codes"] = sorted(set(reasons)) + return result + + +def build_rates_context_observation_v2(snapshot: Mapping, config: Mapping) -> dict: + """Observe explicit collector-known inputs without asserting publication. + + The collector and per-row first-capture/content hashes are caller + declarations. This stateless function does not authenticate timestamps, + verify retained bytes, establish the earliest capture, or enforce stable + first-seen records across calls. First-seen never becomes v1 available_at. + """ + start, end, max_age = _config(config) + snapshot = snapshot if isinstance(snapshot, Mapping) else {} + decision = _forward_time(snapshot.get("decision_at")) + collector = _text(snapshot.get("collector_id")) + series = snapshot.get("series") + top_valid = (set(snapshot) == _FORWARD_ROOT_KEYS and snapshot.get("schema_version") == FORWARD_INPUT_VERSION + and snapshot.get("availability_basis") == _FORWARD_BASIS and collector is not None + and decision is not None and isinstance(series, Mapping) and not set(series).difference(ROLES)) + series = series if isinstance(series, Mapping) else {} + observations = {role: _observe_forward_series(series.get(role), role, start, end, decision, max_age, top_valid) + for role in ROLES} + identities = {} + for role, item in observations.items(): + if item["source_id"] is not None and item["series_id"] is not None: + identities.setdefault((item["source_id"], item["series_id"]), []).append(role) + for repeated_roles in identities.values(): + if len(repeated_roles) > 1: + for role in repeated_roles: + item = observations[role] + item.update(status="unknown", start_percent=None, end_percent=None, change_bp=None) + item["reason_codes"] = sorted(set(item["reason_codes"] + ["ROLE_SOURCE_IDENTITY_AMBIGUOUS"])) + spread = {"status": "unknown", "reason_codes": [], "measurement": "APPROXIMATE_NOMINAL_MINUS_REAL_YIELD_SPREAD", + "source_id": None, "basis": None, "input_series_ids": [], "unit": "percent", + "start_percent": None, "end_percent": None, "change_bp": None} + nominal, real = observations["nominal_10y"], observations["real_10y"] + if any(item["status"] != "declared_available" for item in (nominal, real)): + spread["reason_codes"] = ["COMMON_WINDOW_UNAVAILABLE"] + elif nominal["source_id"] != real["source_id"] or nominal["basis"] != real["basis"]: + spread["reason_codes"] = ["SOURCE_OR_BASIS_MISMATCH"] + else: + first = _rounded(nominal["start_percent"] - real["start_percent"]) + last = _rounded(nominal["end_percent"] - real["end_percent"]) + change = None if first is None or last is None else _rounded(100 * (last - first)) + if change is None: + spread["reason_codes"] = ["SPREAD_NONFINITE"] + else: + spread.update(status="declared_available", source_id=nominal["source_id"], basis=nominal["basis"], + input_series_ids=[nominal["series_id"], real["series_id"]], + start_percent=first, end_percent=last, change_bp=change) + return {"schema_version": FORWARD_OBSERVATION_VERSION, "availability_basis": _FORWARD_BASIS, + "collector_id": collector, "research_status": "UNVALIDATED_RESEARCH", + "assurance": "CALLER_DECLARATIONS_AND_CONSISTENCY_ONLY", "historical_pit_verified": False, + "backtest_eligible": False, "position_control_allowed": False, + "decision_at": _stamp(decision), "window_start": start.isoformat(), "window_end": end.isoformat(), + "window_method": "EXACT_ENDPOINT_DIFFERENCE", "coverage": "ENDPOINTS_ONLY_NOT_SESSION_COVERAGE", + "max_observation_age_days": max_age, "series": observations, "approximate_yield_spread": spread, + "quality": {"status": "declared_available" if all(item["status"] == "declared_available" + for item in observations.values()) and spread["status"] == "declared_available" else "unknown"}} diff --git a/tests/test_rates_context_observer_research.py b/tests/test_rates_context_observer_research.py index 4edb501..8c0334c 100644 --- a/tests/test_rates_context_observer_research.py +++ b/tests/test_rates_context_observer_research.py @@ -333,5 +333,354 @@ def test_unknown_preserves_available_endpoint_metadata(self): self.assertEqual(end["received_at"], "2024-01-04T13:01:00Z") +def forward_fixture(): + base = fixture() + result = {"schema_version": "qsl.rates-context-input.research.v2", + "availability_basis": "collector_first_seen", + "collector_id": "SYNTHETIC_FIXTURE_ONLY:collector-v1", + "decision_at": base["decision_at"], "series": copy.deepcopy(base["series"])} + for role_index, series in enumerate(result["series"].values()): + for row_index, row in enumerate(series["rows"]): + row["source_published_at"] = None + row["first_seen_at"] = row.pop("available_at") + row["capture_sha256"] = f"{100 + role_index:064x}" + row["row_sha256"] = f"{1 + role_index * 10 + row_index:064x}" + return result + + +class ForwardRatesContextObserverTests(unittest.TestCase): + def build(self, snapshot=None, policy=None): + return observer.build_rates_context_observation_v2( + forward_fixture() if snapshot is None else snapshot, config() if policy is None else policy) + + def assert_unknown(self, snapshot, role="nominal_10y"): + result = self.build(snapshot) + self.assertEqual(result["series"][role]["status"], "unknown") + self.assertIsNone(result["series"][role]["change_bp"]) + return result + + def test_explicit_versions_and_collector_identity(self): + result = self.build() + self.assertEqual(result["schema_version"], "qsl.rates-context-observation.research.v2") + self.assertEqual(result["availability_basis"], "collector_first_seen") + self.assertEqual(result["collector_id"], "SYNTHETIC_FIXTURE_ONLY:collector-v1") + self.assertEqual(observer.INPUT_VERSION, "qsl.rates-context-input.research.v1") + self.assertEqual(observer.OBSERVATION_VERSION, "qsl.rates-context-observation.research.v1") + + def test_v1_still_rejects_v2_and_first_seen_fields(self): + result = observer.build_rates_context_observation(forward_fixture(), config()) + self.assertEqual(result["quality"]["status"], "unknown") + self.assertIn("INPUT_INVALID", result["series"]["nominal_10y"]["reason_codes"]) + old = fixture() + old["series"]["nominal_10y"]["rows"][0]["first_seen_at"] = "2024-01-03T13:00:00Z" + result = observer.build_rates_context_observation(old, config()) + self.assertIn("ROW_FIELDS_INVALID", result["series"]["nominal_10y"]["reason_codes"]) + + def test_v2_does_not_accept_v1_or_implicit_version_upgrade(self): + result = self.build(fixture()) + self.assertEqual(result["quality"]["status"], "unknown") + self.assertIn("INPUT_INVALID", result["series"]["nominal_10y"]["reason_codes"]) + + def test_required_root_fields_and_exact_basis(self): + for key, value in (("collector_id", None), ("collector_id", ""), + ("collector_id", " bad "), ("availability_basis", "publisher_time"), + ("schema_version", "unknown"), ("decision_at", "2024-01-04T15:00:00")): + with self.subTest(key=key, value=value): + snapshot = forward_fixture() + snapshot[key] = value + self.assert_unknown(snapshot) + for key in ("collector_id", "availability_basis"): + snapshot = forward_fixture() + del snapshot[key] + self.assert_unknown(snapshot) + + def test_publication_is_explicitly_unknown_and_never_available_at(self): + result = self.build() + end = result["series"]["nominal_10y"]["end_observation"] + self.assertIsNone(end["source_published_at"]) + self.assertEqual(end["first_seen_at"], "2024-01-04T13:00:00Z") + self.assertEqual(end["received_at"], "2024-01-04T13:01:00Z") + self.assertEqual(end["known_at"], end["received_at"]) + def walk(value): + if isinstance(value, dict): + self.assertNotIn("available_at", value) + self.assertNotIn("availability_delay_calendar_days", value) + for nested in value.values(): + walk(nested) + elif isinstance(value, list): + for nested in value: + walk(nested) + walk(result) + + def test_observation_and_consumer_latency_have_distinct_names(self): + nominal = self.build()["series"]["nominal_10y"] + self.assertEqual(nominal["first_seen_delay_calendar_days"], 1) + self.assertEqual(nominal["consumer_receipt_lag_seconds"], 60.0) + self.assertNotIn("receipt_lag_seconds", nominal) + + def test_missing_or_nonnull_publication_is_unknown(self): + for value in ("2024-01-03T13:00:00Z", "", 0, False, object()): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0]["source_published_at"] = value + self.assert_unknown(snapshot) + snapshot = forward_fixture() + del snapshot["series"]["nominal_10y"]["rows"][0]["source_published_at"] + self.assert_unknown(snapshot) + + def test_missing_invalid_clock_or_digest_never_qualifies(self): + for key in ("first_seen_at", "received_at", "revision_id", "capture_sha256", "row_sha256"): + for value in (None, "", "latest", "bad"): + with self.subTest(key=key, value=value): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0][key] = value + if key == "revision_id" and value == "bad": + continue # Opaque immutable revision labels are declarations. + self.assert_unknown(snapshot) + snapshot = forward_fixture() + del snapshot["series"]["nominal_10y"]["rows"][0][key] + self.assert_unknown(snapshot) + + def test_digest_is_lowercase_sha256_shape_not_authenticated_evidence(self): + for key in ("capture_sha256", "row_sha256"): + for value in ("A" * 64, "g" * 64, "a" * 63, "a" * 65, True, 1): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0][key] = value + self.assert_unknown(snapshot) + + def test_receipt_before_first_seen_is_unknown(self): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0]["received_at"] = "2024-01-03T12:59:59Z" + result = self.assert_unknown(snapshot) + self.assertIn("RECEIPT_BEFORE_FIRST_SEEN", result["series"]["nominal_10y"]["reason_codes"]) + self.assertIsNone(result["series"]["nominal_10y"]["start_observation"]["known_at"]) + + def test_first_seen_before_observation_date_is_unknown(self): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0]["first_seen_at"] = "2024-01-01T23:59:59Z" + self.assert_unknown(snapshot) + + def test_exact_arrival_boundary_and_before_first_seen(self): + snapshot = forward_fixture() + for series in snapshot["series"].values(): + for row in series["rows"]: + row["first_seen_at"] = row["received_at"] = "2024-01-04T15:00:00Z" + self.assertEqual(self.build(snapshot)["quality"]["status"], "declared_available") + snapshot["decision_at"] = "2024-01-04T14:59:59.999999Z" + result = self.build(snapshot) + self.assertEqual(result["series"]["nominal_10y"]["visible_observations"], 0) + self.assertEqual(result["quality"]["status"], "unknown") + + def test_old_observation_can_be_known_now_without_backdated_visibility(self): + snapshot = forward_fixture() + snapshot["decision_at"] = "2024-01-10T15:00:00Z" + policy = config() + policy["max_observation_age_days"] = 20 + for series in snapshot["series"].values(): + for row in series["rows"]: + row["first_seen_at"] = row["received_at"] = "2024-01-10T14:00:00Z" + result = self.build(snapshot, policy) + self.assertEqual(result["quality"]["status"], "declared_available") + self.assertFalse(result["historical_pit_verified"]) + snapshot["decision_at"] = "2024-01-09T15:00:00Z" + self.assertEqual(self.build(snapshot, policy)["series"]["nominal_10y"]["visible_observations"], 0) + + def test_today_receipt_never_refreshes_observation_age(self): + snapshot = forward_fixture() + snapshot["decision_at"] = "2024-01-10T15:00:00Z" + for series in snapshot["series"].values(): + for row in series["rows"]: + row["first_seen_at"] = row["received_at"] = snapshot["decision_at"] + result = self.assert_unknown(snapshot) + self.assertIn("OBSERVATION_STALE", result["series"]["nominal_10y"]["reason_codes"]) + + def test_future_or_late_append_preserves_entire_observation(self): + for selector in ("first_seen_at", "received_at", "outside_window"): + snapshot = forward_fixture() + expected = self.build(snapshot) + row = {"observation_date": "2024-01-02", "value": object(), "extra": object()} + if selector == "outside_window": + row["observation_date"] = "2024-01-08" + else: + row[selector] = "2024-01-05T00:00:00Z" + snapshot["series"]["nominal_10y"]["rows"].append(row) + self.assertEqual(self.build(snapshot), expected) + + def test_future_observation_inside_requested_window_is_invisible(self): + snapshot, policy = forward_fixture(), config() + policy["window_end"] = "2024-01-08" + expected = self.build(snapshot, policy) + snapshot["series"]["nominal_10y"]["rows"].append({"observation_date": "2024-01-08", "value": object()}) + self.assertEqual(self.build(snapshot, policy), expected) + + def test_invalid_selector_cannot_claim_invisibility(self): + for key in ("observation_date", "first_seen_at", "received_at"): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0][key] = "bad" + self.assert_unknown(snapshot) + + def test_same_day_two_visible_versions_and_exact_duplicates_are_unknown(self): + for changed in (False, True): + snapshot = forward_fixture() + repeated = copy.deepcopy(snapshot["series"]["nominal_10y"]["rows"][-1]) + if changed: + repeated["revision_id"] = "synthetic-v2" + repeated["value"] = 8.0 + repeated["row_sha256"] = "f" * 64 + snapshot["series"]["nominal_10y"]["rows"].append(repeated) + result = self.assert_unknown(snapshot) + self.assertIn("ROW_ORDER_OR_REVISION_AMBIGUOUS", result["series"]["nominal_10y"]["reason_codes"]) + self.assertIsNone(result["series"]["nominal_10y"]["end_observation"]) + + def test_same_revision_or_content_hash_conflict_is_unknown(self): + snapshot = forward_fixture() + repeated = copy.deepcopy(snapshot["series"]["nominal_10y"]["rows"][-1]) + repeated["value"] = 8.0 + snapshot["series"]["nominal_10y"]["rows"].append(repeated) + result = self.assert_unknown(snapshot) + self.assertIn("ROW_IDENTITY_CONFLICT", result["series"]["nominal_10y"]["reason_codes"]) + snapshot = forward_fixture() + source_rows = snapshot["series"]["nominal_10y"]["rows"] + source_rows[1]["row_sha256"] = source_rows[0]["row_sha256"] + result = self.assert_unknown(snapshot) + self.assertIn("ROW_IDENTITY_CONFLICT", result["series"]["nominal_10y"]["reason_codes"]) + + def test_shared_capture_hash_is_allowed_for_different_rows(self): + result = self.build() + self.assertEqual(result["series"]["nominal_10y"]["status"], "declared_available") + + def test_repeated_call_preserves_declared_first_seen_and_capture_reference(self): + snapshot = forward_fixture() + expected = self.build(snapshot) + self.assertEqual(self.build(copy.deepcopy(snapshot)), expected) + endpoint = expected["series"]["nominal_10y"]["start_observation"] + self.assertEqual(endpoint["capture_sha256"], snapshot["series"]["nominal_10y"]["rows"][0]["capture_sha256"]) + + def test_visible_reverse_order_is_not_repaired(self): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"].reverse() + self.assert_unknown(snapshot) + + def test_partial_evidence_does_not_erase_other_series(self): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0]["first_seen_at"] = None + result = self.assert_unknown(snapshot) + self.assertEqual(result["series"]["real_10y"]["status"], "declared_available") + self.assertEqual(result["approximate_yield_spread"]["status"], "unknown") + + def test_missing_reported_breakeven_keeps_pair_and_approximate_spread(self): + snapshot = forward_fixture() + del snapshot["series"]["breakeven_10y"] + result = self.build(snapshot) + self.assertEqual(result["series"]["nominal_10y"]["change_bp"], 20.0) + self.assertEqual(result["series"]["real_10y"]["change_bp"], 10.0) + self.assertEqual(result["approximate_yield_spread"]["change_bp"], 10.0) + self.assertEqual(result["series"]["breakeven_10y"]["status"], "unknown") + self.assertEqual(result["quality"]["status"], "unknown") + + def test_independent_breakeven_not_replaced_and_negative_yield_valid(self): + result = self.build() + self.assertEqual(result["series"]["breakeven_10y"]["end_percent"], 2.9) + self.assertEqual(result["approximate_yield_spread"]["end_percent"], 3.1) + snapshot = forward_fixture() + for row, value in zip(snapshot["series"]["real_10y"]["rows"], (-1.0, -0.5)): + row["value"] = value + self.assertEqual(self.build(snapshot)["series"]["real_10y"]["change_bp"], 50.0) + + def test_source_basis_unit_and_role_identity_still_apply(self): + for key, value in (("unit", "ratio"), ("basis", "price"), ("source_id", None)): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"][key] = value + self.assert_unknown(snapshot) + snapshot = forward_fixture() + snapshot["series"]["real_10y"]["source_id"] = "other_source" + self.assertEqual(self.build(snapshot)["approximate_yield_spread"]["status"], "unknown") + snapshot = forward_fixture() + snapshot["series"]["real_10y"]["series_id"] = snapshot["series"]["nominal_10y"]["series_id"] + self.assert_unknown(snapshot) + + def test_malformed_rows_and_nonfinite_values_do_not_produce_numbers(self): + for value in (float("nan"), float("inf"), True, "4.0", 10 ** 1000): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0]["value"] = value + result = self.assert_unknown(snapshot) + json.dumps(result, allow_nan=False) + for value in ("not_rows", None, {}): + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"] = value + self.assert_unknown(snapshot) + + def test_extra_ai_or_authority_fields_never_enter_contract(self): + for key in ("ai_narrative", "target_weight", "position_control_allowed", "opportunity"): + snapshot = forward_fixture() + snapshot[key] = True + result = self.assert_unknown(snapshot) + self.assertFalse(result["position_control_allowed"]) + snapshot = forward_fixture() + snapshot["series"]["nominal_10y"]["rows"][0]["available_at"] = "2024-01-02T00:00:00Z" + self.assert_unknown(snapshot) + + def test_timezone_normalization(self): + snapshot = forward_fixture() + row = snapshot["series"]["nominal_10y"]["rows"][-1] + row["first_seen_at"], row["received_at"] = "2024-01-04T08:00:00-05:00", "2024-01-04T08:01:00-05:00" + self.assertEqual(self.build(snapshot), self.build()) + + def test_no_source_or_runtime_authority_and_inputs_unchanged(self): + snapshot, policy = forward_fixture(), config() + before = copy.deepcopy((snapshot, policy)) + result = self.build(snapshot, policy) + self.assertEqual((snapshot, policy), before) + json.dumps(result, allow_nan=False) + self.assertEqual(result["assurance"], "CALLER_DECLARATIONS_AND_CONSISTENCY_ONLY") + self.assertEqual(result["research_status"], "UNVALIDATED_RESEARCH") + for key in ("historical_pit_verified", "backtest_eligible", "position_control_allowed"): + self.assertIs(result[key], False) + + def test_invalid_config_uses_existing_error(self): + policy = config() + policy["max_observation_age_days"] = -1 + with self.assertRaises(observer.ContractError): + self.build(policy=policy) + + + def test_submicrosecond_arrival_cannot_be_truncated_into_visibility(self): + snapshot = forward_fixture() + snapshot["decision_at"] = "2024-01-04T15:00:00Z" + for series in snapshot["series"].values(): + for row in series["rows"]: + row["first_seen_at"] = row["received_at"] = "2024-01-04T15:00:00.0000001Z" + self.assert_unknown(snapshot) + + def test_submicrosecond_clock_reversal_is_not_rounded_equal(self): + snapshot = forward_fixture() + row = snapshot["series"]["nominal_10y"]["rows"][-1] + row["first_seen_at"] = "2024-01-04T13:00:00.0000009Z" + row["received_at"] = "2024-01-04T13:00:00.0000001Z" + self.assert_unknown(snapshot) + + def test_v2_rejects_unsupported_precision_or_normalized_invalid_offsets(self): + for value in ("2024-01-04T15:00:00.1234567Z", "2024-01-04T15:00:00+01:60", + "2024-01-04T15:00:00+00:00:00.0000001", + "2024-01-04T15:00:00.0000000Z", "2024-01-04T15:00:00,1Z"): + for field in ("decision_at", "first_seen_at", "received_at"): + with self.subTest(field=field, value=value): + snapshot = forward_fixture() + if field == "decision_at": + snapshot[field] = value + else: + snapshot["series"]["nominal_10y"]["rows"][-1][field] = value + self.assert_unknown(snapshot) + + def test_microsecond_boundary_remains_exact(self): + snapshot = forward_fixture() + snapshot["decision_at"] = "2024-01-04T15:00:00.000001Z" + for series in snapshot["series"].values(): + for row in series["rows"]: + row["first_seen_at"] = row["received_at"] = snapshot["decision_at"] + self.assertEqual(self.build(snapshot)["quality"]["status"], "declared_available") + snapshot["decision_at"] = "2024-01-04T15:00:00.000000Z" + self.assertEqual(self.build(snapshot)["series"]["nominal_10y"]["visible_observations"], 0) + + if __name__ == "__main__": unittest.main()