The unsafe comparison filters only event time and takes the highest revision. It intentionally ignores availability, tombstones and freshness; it is never exported as a training feature.
+
+
Decision
Source
Feature
Safe value
Safe status
Unsafe value
Uses future knowledge?
SME-A-1
bank
net_inflow
—
no_history
—
False
SME-A-1
tax
invoice_revenue
—
not_available
180000
True
SME-A-1
utility
kwh
1800.0
selected
1800
False
SME-A-2
bank
net_inflow
50000.0
selected
—
True
SME-A-2
tax
invoice_revenue
100000.0
selected
180000
True
SME-A-2
utility
kwh
1800.0
selected
1800
False
SME-A-3
bank
net_inflow
—
deleted
—
False
SME-A-3
tax
invoice_revenue
100000.0
selected
180000
True
SME-A-3
utility
kwh
1800.0
selected
1800
False
SME-A-4
bank
net_inflow
—
deleted
—
False
SME-A-4
tax
invoice_revenue
180000.0
selected
220000
True
SME-A-4
utility
kwh
1800.0
selected
1800
False
SME-A-5
bank
net_inflow
—
deleted
—
False
SME-A-5
tax
invoice_revenue
220000.0
selected
220000
False
SME-A-5
utility
kwh
—
stale
1800
False
SME-B-1
bank
net_inflow
—
no_history
—
False
SME-B-1
tax
invoice_revenue
—
not_available
90000
True
SME-B-1
utility
kwh
—
stale
1300
True
SME-B-2
bank
net_inflow
30000.0
selected
30000
False
SME-B-2
tax
invoice_revenue
—
not_available
90000
True
SME-B-2
utility
kwh
—
stale
1300
True
SME-B-3
bank
net_inflow
30000.0
selected
30000
False
SME-B-3
tax
invoice_revenue
—
not_available
90000
True
SME-B-3
utility
kwh
—
stale
1300
True
SME-B-4
bank
net_inflow
30000.0
selected
30000
False
SME-B-4
tax
invoice_revenue
90000.0
selected
90000
False
SME-B-4
utility
kwh
1300.0
selected
1300
False
SME-B-5
bank
net_inflow
30000.0
selected
30000
False
SME-B-5
tax
invoice_revenue
90000.0
selected
90000
False
SME-B-5
utility
kwh
—
stale
1300
False
SME-C-1
bank
net_inflow
—
no_history
—
False
SME-C-1
tax
invoice_revenue
—
no_history
—
False
SME-C-1
utility
kwh
—
no_history
—
False
SME-C-2
bank
net_inflow
—
no_history
—
False
SME-C-2
tax
invoice_revenue
—
no_history
—
False
SME-C-2
utility
kwh
—
no_history
—
False
SME-C-3
bank
net_inflow
—
no_history
—
False
SME-C-3
tax
invoice_revenue
—
no_history
—
False
SME-C-3
utility
kwh
—
no_history
—
False
SME-C-4
bank
net_inflow
—
no_history
—
False
SME-C-4
tax
invoice_revenue
—
no_history
—
False
SME-C-4
utility
kwh
—
no_history
—
False
SME-C-5
bank
net_inflow
—
no_history
—
False
SME-C-5
tax
invoice_revenue
—
no_history
—
False
SME-C-5
utility
kwh
—
no_history
—
False
+
Trace each selected observation
Knowledge time = max(publication, ingestion). Every selected observation was available at or before the decision. A later historical correction cannot overwrite an earlier snapshot.
+
Decision
Source
Feature
Record
Revision
Event time
Available time
SME-A-1
utility
kwh
power-a
1
2026-02-01T00:00:00.000000+00:00
2026-02-03T00:00:00.000000+00:00
SME-A-2
bank
net_inflow
cash-a
1
2026-02-10T00:00:00.000000+00:00
2026-02-11T00:00:00.000000+00:00
SME-A-2
tax
invoice_revenue
tax-a-v1
1
2026-01-31T00:00:00.000000+00:00
2026-02-06T00:00:00.000000+00:00
SME-A-2
utility
kwh
power-a
1
2026-02-01T00:00:00.000000+00:00
2026-02-03T00:00:00.000000+00:00
SME-A-3
tax
invoice_revenue
tax-a-v1
1
2026-01-31T00:00:00.000000+00:00
2026-02-06T00:00:00.000000+00:00
SME-A-3
utility
kwh
power-a
1
2026-02-01T00:00:00.000000+00:00
2026-02-03T00:00:00.000000+00:00
SME-A-4
tax
invoice_revenue
tax-a-v2
2
2026-01-31T00:00:00.000000+00:00
2026-03-11T00:00:00.000000+00:00
SME-A-4
utility
kwh
power-a
1
2026-02-01T00:00:00.000000+00:00
2026-02-03T00:00:00.000000+00:00
SME-A-5
tax
invoice_revenue
tax-a-feb
1
2026-02-28T00:00:00.000000+00:00
2026-03-20T00:00:00.000000+00:00
SME-B-2
bank
net_inflow
cash-b-v2
2
2026-02-10T00:00:00.000000+00:00
2026-02-13T00:00:00.000000+00:00
SME-B-3
bank
net_inflow
cash-b-v2
2
2026-02-10T00:00:00.000000+00:00
2026-02-13T00:00:00.000000+00:00
SME-B-4
bank
net_inflow
cash-b-v2
2
2026-02-10T00:00:00.000000+00:00
2026-02-13T00:00:00.000000+00:00
SME-B-4
tax
invoice_revenue
tax-b
1
2026-01-31T00:00:00.000000+00:00
2026-03-01T00:00:00.000000+00:00
SME-B-4
utility
kwh
power-b
1
2026-02-01T00:00:00.000000+00:00
2026-03-02T00:00:00.000000+00:00
SME-B-5
bank
net_inflow
cash-b-v2
2
2026-02-10T00:00:00.000000+00:00
2026-02-13T00:00:00.000000+00:00
SME-B-5
tax
invoice_revenue
tax-b
1
2026-01-31T00:00:00.000000+00:00
2026-03-01T00:00:00.000000+00:00
+
Replay locally
pitbridge verify --out demo checks file hashes, recomputes every snapshot through both implementations, and compares all derived evidence. Hashes detect accidental edits; they are not a digital signature.
+
+
diff --git a/demo/snapshots.csv b/demo/snapshots.csv
new file mode 100644
index 0000000..e7a613b
--- /dev/null
+++ b/demo/snapshots.csv
@@ -0,0 +1,46 @@
+decision_id,entity_id,decision_at,source,feature,value,status,record_id,revision,event_at,published_at,ingested_at,available_at
+SME-A-1,SME-A,2026-02-04T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-A-1,SME-A,2026-02-04T00:00:00.000000+00:00,tax,invoice_revenue,,not_available,,,,,,
+SME-A-1,SME-A,2026-02-04T00:00:00.000000+00:00,utility,kwh,1800.0,selected,power-a,1,2026-02-01T00:00:00.000000+00:00,2026-02-02T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00
+SME-A-2,SME-A,2026-02-15T00:00:00.000000+00:00,bank,net_inflow,50000.0,selected,cash-a,1,2026-02-10T00:00:00.000000+00:00,2026-02-10T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00
+SME-A-2,SME-A,2026-02-15T00:00:00.000000+00:00,tax,invoice_revenue,100000.0,selected,tax-a-v1,1,2026-01-31T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-06T00:00:00.000000+00:00,2026-02-06T00:00:00.000000+00:00
+SME-A-2,SME-A,2026-02-15T00:00:00.000000+00:00,utility,kwh,1800.0,selected,power-a,1,2026-02-01T00:00:00.000000+00:00,2026-02-02T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00
+SME-A-3,SME-A,2026-02-20T00:00:00.000000+00:00,bank,net_inflow,,deleted,,,,,,
+SME-A-3,SME-A,2026-02-20T00:00:00.000000+00:00,tax,invoice_revenue,100000.0,selected,tax-a-v1,1,2026-01-31T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-06T00:00:00.000000+00:00,2026-02-06T00:00:00.000000+00:00
+SME-A-3,SME-A,2026-02-20T00:00:00.000000+00:00,utility,kwh,1800.0,selected,power-a,1,2026-02-01T00:00:00.000000+00:00,2026-02-02T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00
+SME-A-4,SME-A,2026-03-15T00:00:00.000000+00:00,bank,net_inflow,,deleted,,,,,,
+SME-A-4,SME-A,2026-03-15T00:00:00.000000+00:00,tax,invoice_revenue,180000.0,selected,tax-a-v2,2,2026-01-31T00:00:00.000000+00:00,2026-03-10T00:00:00.000000+00:00,2026-03-11T00:00:00.000000+00:00,2026-03-11T00:00:00.000000+00:00
+SME-A-4,SME-A,2026-03-15T00:00:00.000000+00:00,utility,kwh,1800.0,selected,power-a,1,2026-02-01T00:00:00.000000+00:00,2026-02-02T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00
+SME-A-5,SME-A,2026-03-25T00:00:00.000000+00:00,bank,net_inflow,,deleted,,,,,,
+SME-A-5,SME-A,2026-03-25T00:00:00.000000+00:00,tax,invoice_revenue,220000.0,selected,tax-a-feb,1,2026-02-28T00:00:00.000000+00:00,2026-03-05T00:00:00.000000+00:00,2026-03-20T00:00:00.000000+00:00,2026-03-20T00:00:00.000000+00:00
+SME-A-5,SME-A,2026-03-25T00:00:00.000000+00:00,utility,kwh,,stale,,,,,,
+SME-B-1,SME-B,2026-02-04T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-B-1,SME-B,2026-02-04T00:00:00.000000+00:00,tax,invoice_revenue,,not_available,,,,,,
+SME-B-1,SME-B,2026-02-04T00:00:00.000000+00:00,utility,kwh,,stale,,,,,,
+SME-B-2,SME-B,2026-02-15T00:00:00.000000+00:00,bank,net_inflow,30000.0,selected,cash-b-v2,2,2026-02-10T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00
+SME-B-2,SME-B,2026-02-15T00:00:00.000000+00:00,tax,invoice_revenue,,not_available,,,,,,
+SME-B-2,SME-B,2026-02-15T00:00:00.000000+00:00,utility,kwh,,stale,,,,,,
+SME-B-3,SME-B,2026-02-20T00:00:00.000000+00:00,bank,net_inflow,30000.0,selected,cash-b-v2,2,2026-02-10T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00
+SME-B-3,SME-B,2026-02-20T00:00:00.000000+00:00,tax,invoice_revenue,,not_available,,,,,,
+SME-B-3,SME-B,2026-02-20T00:00:00.000000+00:00,utility,kwh,,stale,,,,,,
+SME-B-4,SME-B,2026-03-15T00:00:00.000000+00:00,bank,net_inflow,30000.0,selected,cash-b-v2,2,2026-02-10T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00
+SME-B-4,SME-B,2026-03-15T00:00:00.000000+00:00,tax,invoice_revenue,90000.0,selected,tax-b,1,2026-01-31T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-03-01T00:00:00.000000+00:00,2026-03-01T00:00:00.000000+00:00
+SME-B-4,SME-B,2026-03-15T00:00:00.000000+00:00,utility,kwh,1300.0,selected,power-b,1,2026-02-01T00:00:00.000000+00:00,2026-02-03T00:00:00.000000+00:00,2026-03-02T00:00:00.000000+00:00,2026-03-02T00:00:00.000000+00:00
+SME-B-5,SME-B,2026-03-25T00:00:00.000000+00:00,bank,net_inflow,30000.0,selected,cash-b-v2,2,2026-02-10T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00,2026-02-13T00:00:00.000000+00:00
+SME-B-5,SME-B,2026-03-25T00:00:00.000000+00:00,tax,invoice_revenue,90000.0,selected,tax-b,1,2026-01-31T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-03-01T00:00:00.000000+00:00,2026-03-01T00:00:00.000000+00:00
+SME-B-5,SME-B,2026-03-25T00:00:00.000000+00:00,utility,kwh,,stale,,,,,,
+SME-C-1,SME-C,2026-02-04T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-C-1,SME-C,2026-02-04T00:00:00.000000+00:00,tax,invoice_revenue,,no_history,,,,,,
+SME-C-1,SME-C,2026-02-04T00:00:00.000000+00:00,utility,kwh,,no_history,,,,,,
+SME-C-2,SME-C,2026-02-15T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-C-2,SME-C,2026-02-15T00:00:00.000000+00:00,tax,invoice_revenue,,no_history,,,,,,
+SME-C-2,SME-C,2026-02-15T00:00:00.000000+00:00,utility,kwh,,no_history,,,,,,
+SME-C-3,SME-C,2026-02-20T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-C-3,SME-C,2026-02-20T00:00:00.000000+00:00,tax,invoice_revenue,,no_history,,,,,,
+SME-C-3,SME-C,2026-02-20T00:00:00.000000+00:00,utility,kwh,,no_history,,,,,,
+SME-C-4,SME-C,2026-03-15T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-C-4,SME-C,2026-03-15T00:00:00.000000+00:00,tax,invoice_revenue,,no_history,,,,,,
+SME-C-4,SME-C,2026-03-15T00:00:00.000000+00:00,utility,kwh,,no_history,,,,,,
+SME-C-5,SME-C,2026-03-25T00:00:00.000000+00:00,bank,net_inflow,,no_history,,,,,,
+SME-C-5,SME-C,2026-03-25T00:00:00.000000+00:00,tax,invoice_revenue,,no_history,,,,,,
+SME-C-5,SME-C,2026-03-25T00:00:00.000000+00:00,utility,kwh,,no_history,,,,,,
diff --git a/demo/summary.json b/demo/summary.json
new file mode 100644
index 0000000..7ab0889
--- /dev/null
+++ b/demo/summary.json
@@ -0,0 +1,17 @@
+{
+ "data_kind": "synthetic",
+ "decisions": 15,
+ "different_selections": 16,
+ "future_knowledge_rows": 11,
+ "independent_replay": true,
+ "observations": 11,
+ "schema_version": 1,
+ "snapshot_rows": 45,
+ "status_counts": {
+ "deleted": 3,
+ "no_history": 17,
+ "not_available": 4,
+ "selected": 16,
+ "stale": 5
+ }
+}
diff --git a/docs/DATA_CONTRACT.md b/docs/DATA_CONTRACT.md
new file mode 100644
index 0000000..16d4434
--- /dev/null
+++ b/docs/DATA_CONTRACT.md
@@ -0,0 +1,45 @@
+# Input and output contract
+
+The JSON root requires `observations`, `decisions` and `specs` arrays. Optional `data_kind` is `synthetic` or `user_supplied` (default). Extra fields, duplicate JSON keys and non-finite JSON constants are rejected. The provenance label is supplied by the caller; it is not externally authenticated.
+
+## Observations
+
+| Field | Type / rule |
+| :--- | :--- |
+| `record_id` | Unique, nonempty trimmed string. |
+| `entity_id`, `source`, `feature` | Nonempty trimmed strings; each source/feature must have a spec. |
+| `event_at` | Timezone-aware ISO-8601 event instant. |
+| `published_at` | Timezone-aware instant, no earlier than event time. |
+| `ingested_at` | Timezone-aware instant, no earlier than event time. |
+| `revision` | Integer in `[1, 2^63-1]`; booleans are rejected. |
+| `value` | Finite number for an active record; null for a tombstone. |
+| `deleted` | Boolean, default false. |
+
+The logical version key `(entity_id, source, feature, event_at, revision)` must also be unique. Value units are governed upstream: the demo uses CNY for invoice revenue and bank net inflow, kWh for utility consumption. Unit conversion and numeric aggregation are outside this implementation.
+
+## Decisions and specs
+
+| Section | Fields |
+| :--- | :--- |
+| `decisions` | Unique `decision_id`, `entity_id`, timezone-aware `decision_at`. |
+| `specs` | Unique pair `source`, `feature`; optional `max_age_days`, null or within `[0,3652500]`. |
+
+At least one decision and one spec are required; observations may be empty. The snapshot has exactly `len(decisions) * len(specs)` rows and a stable decision/source/feature order.
+
+## Example
+
+```json
+{
+ "data_kind": "synthetic",
+ "observations": [{
+ "record_id": "tax-001-v1", "entity_id": "SME-001", "source": "tax",
+ "feature": "invoice_revenue", "event_at": "2026-01-31T00:00:00Z",
+ "published_at": "2026-02-05T00:00:00Z", "ingested_at": "2026-02-06T00:00:00Z",
+ "revision": 1, "value": 100000, "deleted": false
+ }],
+ "decisions": [{"decision_id": "app-001", "entity_id": "SME-001", "decision_at": "2026-02-15T00:00:00Z"}],
+ "specs": [{"source": "tax", "feature": "invoice_revenue", "max_age_days": 90}]
+}
+```
+
+`snapshots.csv` exports the decision identity/time, source, feature, value, status, record ID, revision, event/publication/ingestion/availability times. Missing lineage fields are empty CSV cells. `inputs.json` retains the full normalized history, including revisions not chosen for a decision.
diff --git a/docs/INTERVIEW.md b/docs/INTERVIEW.md
new file mode 100644
index 0000000..9d5b97a
--- /dev/null
+++ b/docs/INTERVIEW.md
@@ -0,0 +1,31 @@
+# 中文面试说明
+
+## 30 秒介绍
+
+我做的是金融时点数据组件。历史数据常常迟到、修订或者被撤销,只按照业务发生日期关联,会让模型看到当时银行还拿不到的信息。我同时约束发生时间和可用时间,先选择当时已知的修订版本,再处理撤销与过期;使用 SQL 生成特征,用独立 Python 实现重放,并导出来源记录级证据。
+
+## 可以使用的简历表述
+
+> 构建金融多源数据时点特征组件,使用 SQLite 窗口函数处理公开/入库延迟、历史修订、撤销及特征有效期,保留记录级溯源;用独立 Python 枚举及随机历史对照验证 SQL 输出,在 45 次合成特征查询中识别事件时间基线的 11 次未来信息使用,并实现报告哈希校验与语义重放。
+
+这是候选表述。要能解释代码和假设;示例是合成数据,不能描述为银行线上项目。
+
+## 常见追问
+
+**为什么 `event_at` 不够?** 一月份的数据可能三月份才入库。二月份做决策时,即使数据“属于一月份”,仍然不可用。
+
+**为什么可用时间取两个时间的最大值?** 当前合同要求数据已经公开,而且系统已经拿到,两个条件都满足才可用。若真实系统的可得性规则不同,应修改源合同,不能凭空猜测。
+
+**为什么不按照最后入库时间选版本?** 旧版本可能迟到。最高已知源版本才是该观察时点的权威版本;乱序入库不应让旧版本覆盖新版本。
+
+**撤销为什么要先选版本再过滤?** 先删除撤销行再排名,会让旧版本重新变成可用,等于把已撤销的数据复活。
+
+**11 和 16 分别代表什么?** 11 次基线选择使用了未来可用信息。16 次最终选择不同,包含撤销和过期的影响;不能把 16 全部说成未来泄漏。
+
+**哈希已经匹配,为什么还要重放?** 哈希只能证明文件和记录的摘要一致。有人改了指标再更新哈希,文件仍可能逻辑错误,因此要从输入重算。
+
+**能支持多少数据?** 当前没有分布式性能验证。面试可以说明 SQL 逻辑、索引以及如何分区或优化候选连接,不应编造吞吐量。
+
+## 建议展示路径
+
+先打开 `demo/comparison.csv` 看 `SME-A-2 / tax` 的未来修订;再打开 `demo/snapshots.csv` 看选中的 `tax-a-v1`;最后展示 `test_historical_revision_cannot_change_old_decision` 和 `test_rehashed_false_metric_still_fails_semantic_replay`。
diff --git a/docs/METHODOLOGY.md b/docs/METHODOLOGY.md
new file mode 100644
index 0000000..dbbcc28
--- /dev/null
+++ b/docs/METHODOLOGY.md
@@ -0,0 +1,44 @@
+# Methodology
+
+## Two clocks and a source revision
+
+An observation has event time `e`, publication time `p`, ingestion time `i`, and availability `a = max(p, i)`. This assumes both public release and ingestion are necessary for the decision system to use the record. It also supports pre-ingested embargoed records: ingestion alone does not make them eligible.
+
+For decision time `t`, the eligible set requires `e <= t` and `a <= t`. All instants are normalized to UTC at microsecond precision. Equality is inclusive. The contract rejects naive timestamps and publication/ingestion preceding the represented observation event.
+
+The logical observation key is `(entity_id, source, feature, event_at)`. Within the eligible set, choose the largest positive source `revision`. A higher version remains authoritative even if a lower version was delivered later. Duplicate logical key/revision pairs are rejected rather than resolved by incidental row order.
+
+A tombstone is a revision with `deleted=true` and `value=null`. Version resolution happens **before** removing tombstones, so a revoked observation cannot fall back to its older revision. It does not revoke the entire feature series: an earlier, separate active observation may still be selected.
+
+For an optional freshness limit `w`, use `0 <= t-e <= w` after resolving versions. The latest active event within the window is selected. Days mean exactly 86,400 seconds, rounded to microseconds; there is no business-calendar adjustment.
+
+## Missingness precedence
+
+| Status | Meaning |
+| :--- | :--- |
+| `no_history` | No observation for this entity/source/feature has event time at or before the decision. |
+| `not_available` | Historical events exist, but none were available by the decision. |
+| `deleted` | Known observations exist; their highest known versions are all tombstones. |
+| `stale` | Known active versions exist, but none pass the freshness window. |
+| `selected` | A valid record and its lineage were selected. |
+
+If known tombstones coexist with old active observations, and those active observations are stale, the status is `stale`. Missingness is not imputed here.
+
+## An auditable counterexample
+
+`tax-a-v1` represents 100,000 CNY of January invoice revenue and is available February 6. Revision 2 changes the value to 180,000 CNY but is not available until March 11. The February 15 decision must keep version 1. The fixture separately tests March's late ingestion, a revoked bank cash observation, UTC offset equivalence, and delayed delivery of an older revision.
+
+The unsafe baseline only constrains event time, choosing the newest event/highest revision. It intentionally ignores availability, deletion and freshness. `future_knowledge_rows` counts unavailable baseline records; `different_selections` also counts non-leakage contract differences. The fixture is constructed to expose these failures and is not a representative sample of bank data.
+
+## Verification boundaries
+
+SQLite resolves eligible versions and event ranks using SQL window functions. The independent Python oracle enumerates histories and versions without calling SQL. Tests compare both on randomized data, shuffled input order and appended future revisions.
+
+`verify` first checks artifact hashes, then derives all outputs again from the normalized input snapshot and compares them byte-for-byte. A changed metric plus an updated hash fails replay. An attacker who rewrites the source inputs and regenerates the whole bundle can produce a consistent new bundle: this workflow is not authentication or immutable source storage.
+
+## Sources and design context
+
+- [Feast point-in-time joins](https://docs.feast.dev/getting-started/concepts/point-in-time-joins): historical feature retrieval and freshness windows. PITBridge explicitly uses the maximum of publication and ingestion time in addition to event time.
+- [SQLite window functions](https://www.sqlite.org/windowfunctions.html): ranking for observation revisions and event selection.
+
+No Feast code was copied and PITBridge is not a Feast adapter. There are no performance or production-readiness claims.
diff --git a/docs/evidence.png b/docs/evidence.png
new file mode 100644
index 0000000..99e68bd
Binary files /dev/null and b/docs/evidence.png differ
diff --git a/pyproject.toml b/pyproject.toml
new file mode 100644
index 0000000..b50025e
--- /dev/null
+++ b/pyproject.toml
@@ -0,0 +1,20 @@
+[build-system]
+requires = ["setuptools>=68"]
+build-backend = "setuptools.build_meta"
+
+[project]
+name = "pitbridge"
+version = "0.1.0"
+description = "Availability-aware financial feature snapshots with revision lineage"
+requires-python = ">=3.11"
+license = {text = "MIT"}
+dependencies = []
+
+[project.optional-dependencies]
+figures = ["matplotlib>=3.8,<4"]
+
+[project.scripts]
+pitbridge = "pitbridge.cli:main"
+
+[tool.setuptools.packages.find]
+where = ["src"]
diff --git a/scripts/render_figures.py b/scripts/render_figures.py
new file mode 100644
index 0000000..4c31e81
--- /dev/null
+++ b/scripts/render_figures.py
@@ -0,0 +1,58 @@
+"""Render a source-linked figure from the saved evidence (optional Matplotlib)."""
+
+import csv
+from datetime import datetime
+import json
+from pathlib import Path
+
+import matplotlib
+matplotlib.use("Agg")
+import matplotlib.dates as mdates
+import matplotlib.pyplot as plt
+
+
+ROOT = Path(__file__).resolve().parents[1]
+plt.rcParams.update({"font.family":"DejaVu Sans", "font.size":11, "text.color":"#e6edf7", "axes.labelcolor":"#a9bad0", "xtick.color":"#a9bad0", "ytick.color":"#a9bad0", "axes.edgecolor":"#314460", "axes.facecolor":"#101d30", "figure.facecolor":"#0c1422"})
+
+
+def main():
+ inputs=json.loads((ROOT/"demo/inputs.json").read_text())
+ summary=json.loads((ROOT/"demo/summary.json").read_text())
+ comparison=list(csv.DictReader((ROOT/"demo/comparison.csv").open()))
+ row=next(r for r in comparison if r["decision_id"]=="SME-A-2" and r["source"]=="tax")
+ records=[r for r in inputs["observations"] if r["record_id"] in ("tax-a-v1","tax-a-v2")]
+ fig,(timeline,bars)=plt.subplots(2,1,figsize=(13,7.5),gridspec_kw={"height_ratios":[1.1,1]},layout="constrained")
+ fig.suptitle("PITBridge | What was known at decision time?",fontsize=23,fontweight="bold",x=.055,ha="left")
+ cutoff=datetime.fromisoformat("2026-02-15T00:00:00+00:00")
+ for y,record in zip((1,0),records):
+ event=datetime.fromisoformat(record["event_at"])
+ available=datetime.fromisoformat(max(record["published_at"],record["ingested_at"]))
+ color="#64e5c4" if available<=cutoff else "#efad70"
+ timeline.plot([event,available],[y,y],color=color,linewidth=3,marker="o",markersize=8)
+ timeline.annotate(f"Available {available:%b %d}",(available,y),xytext=(0,13),textcoords="offset points",color=color,ha="center")
+ timeline.axvline(cutoff,color="#97c6ff",linestyle="--",linewidth=1.7)
+ timeline.annotate("Decision: Feb 15",(cutoff,.5),xytext=(8,0),textcoords="offset points",color="#97c6ff")
+ timeline.set_yticks([0,1],["Revision 2: 180,000 CNY","Revision 1: 100,000 CNY"])
+ timeline.set_ylim(-.5,1.55)
+ timeline.xaxis.set_major_formatter(mdates.DateFormatter("%b %d"))
+ timeline.xaxis.set_major_locator(mdates.WeekdayLocator(interval=1))
+ timeline.set_title("The event date is the same. The knowledge date is different.",loc="left",fontsize=13,pad=12)
+ values=[float(row["safe_value"]),float(row["unsafe_value"])]
+ bars.barh([1,0],values,color=["#64e5c4","#efad70"],height=.5)
+ bars.set_yticks([1,0],["Safe snapshot","Unsafe event-only join"])
+ bars.set_xlim(0,max(values)*1.25)
+ bars.set_xlabel("Invoice revenue / CNY")
+ bars.set_title("SME-A / Feb 15: a later correction cannot enter this decision",loc="left",fontsize=13,pad=12)
+ for y,value in zip((1,0),values):
+ bars.text(value+5000,y,f"{value:,.0f}",va="center",fontweight="bold")
+ for ax in (timeline,bars):
+ ax.spines[["top","right"]].set_visible(False)
+ ax.grid(axis="x",color="#536981",alpha=.18)
+ ax.set_axisbelow(True)
+ fig.supxlabel(f"Synthetic fixture · {summary['snapshot_rows']} lookups · {summary['future_knowledge_rows']} future-data baseline selections · Source: demo/inputs.json + comparison.csv",fontsize=10,color="#8ca6bf")
+ fig.savefig(ROOT/"docs/evidence.png",dpi=150)
+ plt.close(fig)
+
+
+if __name__=="__main__":
+ main()
diff --git a/src/pitbridge/__init__.py b/src/pitbridge/__init__.py
new file mode 100644
index 0000000..8e6fdcf
--- /dev/null
+++ b/src/pitbridge/__init__.py
@@ -0,0 +1,6 @@
+"""PITBridge: versioned features as they were known at decision time."""
+
+from .core import Decision, FeatureSpec, Observation, Snapshot, build_snapshot, reference_snapshot
+
+__version__ = "0.1.0"
+__all__ = ["Decision", "FeatureSpec", "Observation", "Snapshot", "build_snapshot", "reference_snapshot"]
diff --git a/src/pitbridge/__main__.py b/src/pitbridge/__main__.py
new file mode 100644
index 0000000..eb53e2f
--- /dev/null
+++ b/src/pitbridge/__main__.py
@@ -0,0 +1,3 @@
+from .cli import main
+
+raise SystemExit(main())
diff --git a/src/pitbridge/bundle.py b/src/pitbridge/bundle.py
new file mode 100644
index 0000000..5995c1c
--- /dev/null
+++ b/src/pitbridge/bundle.py
@@ -0,0 +1,138 @@
+"""Replayable JSON/CSV evidence and a standalone browser report."""
+
+from collections import Counter
+import csv
+from dataclasses import asdict
+from hashlib import sha256
+from html import escape
+import json
+from pathlib import Path
+
+from .core import Decision, FeatureSpec, Observation, build_snapshot, reference_snapshot, validate
+
+
+FILES = ("inputs.json", "snapshots.csv", "comparison.csv", "summary.json", "report.html")
+
+
+def canonical(value) -> str:
+ return json.dumps(value, ensure_ascii=False, sort_keys=True, indent=2, allow_nan=False) + "\n"
+
+
+def read_json(path):
+ def bad(value):
+ raise ValueError(f"non-finite JSON constant: {value}")
+ def pairs(items):
+ result = {}
+ for key, value in items:
+ if key in result:
+ raise ValueError(f"duplicate JSON key: {key}")
+ result[key] = value
+ return result
+ return json.loads(Path(path).read_text(encoding="utf-8"), parse_constant=bad, object_pairs_hook=pairs)
+
+
+def decode_inputs(data):
+ if not isinstance(data, dict) or not {"observations", "decisions", "specs"} <= set(data) or set(data) - {"observations", "decisions", "specs", "data_kind"}:
+ raise ValueError("input object requires observations, decisions, specs; only data_kind is optional")
+ if data.get("data_kind", "user_supplied") not in ("synthetic", "user_supplied"):
+ raise ValueError("data_kind must be synthetic or user_supplied")
+ if any(not isinstance(data[key], list) for key in ("observations", "decisions", "specs")):
+ raise ValueError("all input sections must be arrays")
+ try:
+ return validate([Observation(**row) for row in data["observations"]], [Decision(**row) for row in data["decisions"]], [FeatureSpec(**row) for row in data["specs"]])
+ except TypeError as exc:
+ raise ValueError(f"input schema mismatch: {exc}") from exc
+
+
+def csv_text(rows, fields):
+ import io
+ stream = io.StringIO(newline="")
+ writer = csv.DictWriter(stream, fieldnames=fields, lineterminator="\n")
+ writer.writeheader()
+ writer.writerows(rows)
+ return stream.getvalue()
+
+
+def compare_unsafe(obs, rows):
+ """A deliberately unsafe event-only baseline, clearly isolated from features."""
+ comparison = []
+ for row in rows:
+ candidates = [o for o in obs if (o.entity_id, o.source, o.feature) == (row.entity_id, row.source, row.feature) and o.event_at <= row.decision_at]
+ baseline = max(candidates, key=lambda o: (o.event_at, o.revision)) if candidates else None
+ comparison.append({
+ "decision_id": row.decision_id, "source": row.source, "feature": row.feature,
+ "safe_record_id": row.record_id, "safe_value": row.value, "safe_status": row.status,
+ "unsafe_record_id": baseline.record_id if baseline else None,
+ "unsafe_value": baseline.value if baseline else None,
+ "unsafe_uses_future_knowledge": bool(baseline and baseline.available_at > row.decision_at),
+ "different_selection": row.record_id != (baseline.record_id if baseline else None),
+ })
+ return comparison
+
+
+def render_report(summary, comparison, snapshots):
+ kind = "SYNTHETIC FINANCIAL DATA" if summary["data_kind"] == "synthetic" else "USER-SUPPLIED DATA"
+ provenance = "Synthetic demonstration · No real borrower data" if summary["data_kind"] == "synthetic" else "User-supplied inputs · Source authenticity has not been independently audited"
+ counts = " · ".join(f"{escape(k)}: {v}" for k, v in summary["status_counts"].items())
+ comparison_rows = "".join(
+ '
' +
+ "".join(f"
{escape(str(r[k] if r[k] is not None else '—'))}
" for k in ("decision_id", "source", "feature", "safe_value", "safe_status", "unsafe_value", "unsafe_uses_future_knowledge")) + "
"
+ for r in comparison
+ )
+ lineage_rows = "".join("
" + "".join(f"
{escape(str(getattr(r,k) if getattr(r,k) is not None else '—'))}
" for k in ("decision_id", "source", "feature", "record_id", "revision", "event_at", "available_at")) + "
" for r in snapshots if r.status == "selected")
+ return f"""
+
+PITBridge · Decision-time evidence
+
PITBRIDGE / {kind}
+
What was known when the decision was made?
+
Availability-aware SQL snapshots with publication, ingestion, revision and record-level lineage.
Independent replay: SQLite output matches the Python enumeration oracle. {counts}
+
Inspect the leakage counterexample
The unsafe comparison filters only event time and takes the highest revision. It intentionally ignores availability, tombstones and freshness; it is never exported as a training feature.
+
+
Decision
Source
Feature
Safe value
Safe status
Unsafe value
Uses future knowledge?
{comparison_rows}
+
Trace each selected observation
Knowledge time = max(publication, ingestion). Every selected observation was available at or before the decision. A later historical correction cannot overwrite an earlier snapshot.
+
Decision
Source
Feature
Record
Revision
Event time
Available time
{lineage_rows}
+
Replay locally
pitbridge verify --out demo checks file hashes, recomputes every snapshot through both implementations, and compares all derived evidence. Hashes detect accidental edits; they are not a digital signature.
+
+\n"""
+
+
+def artifacts(data):
+ obs, dec, specs = decode_inputs(data)
+ normalized = {"data_kind": data.get("data_kind", "user_supplied"), "observations": [asdict(r) for r in obs], "decisions": [asdict(r) for r in dec], "specs": [asdict(r) for r in specs]}
+ rows = build_snapshot(obs, dec, specs)
+ if rows != reference_snapshot(obs, dec, specs):
+ raise ValueError("SQL/reference replay mismatch")
+ comparison = compare_unsafe(obs, rows)
+ summary = {"schema_version": 1, "data_kind": normalized["data_kind"], "observations": len(obs), "decisions": len(dec), "snapshot_rows": len(rows), "status_counts": dict(sorted(Counter(r.status for r in rows).items())), "future_knowledge_rows": sum(r["unsafe_uses_future_knowledge"] for r in comparison), "different_selections": sum(r["different_selection"] for r in comparison), "independent_replay": True}
+ # All outputs are derived from normalized inputs; there is no wall-clock timestamp.
+ return {"inputs.json": canonical(normalized), "snapshots.csv": csv_text([asdict(r) for r in rows], list(asdict(rows[0]))), "comparison.csv": csv_text(comparison, list(comparison[0])), "summary.json": canonical(summary), "report.html": render_report(summary, comparison, rows)}
+
+
+def write_bundle(data, out):
+ files = artifacts(data)
+ target = Path(out)
+ target.mkdir(parents=True, exist_ok=True)
+ for name, content in files.items():
+ (target / name).write_text(content, encoding="utf-8", newline="")
+ manifest = {"schema_version": 1, "engine": "pitbridge/0.1.0", "files": {name: sha256(content.encode()).hexdigest() for name, content in sorted(files.items())}}
+ (target / "manifest.json").write_text(canonical(manifest), encoding="utf-8")
+ return read_json(target / "summary.json")
+
+
+def verify_bundle(out):
+ target = Path(out)
+ manifest = read_json(target / "manifest.json")
+ if not isinstance(manifest, dict) or not isinstance(manifest.get("files"), dict) or manifest.get("schema_version") != 1 or manifest.get("engine") != "pitbridge/0.1.0" or set(manifest["files"]) != set(FILES):
+ raise ValueError("unsupported manifest or unexpected artifact paths")
+ for name in FILES:
+ if sha256((target / name).read_bytes()).hexdigest() != manifest["files"][name]:
+ raise ValueError(f"hash mismatch: {name}")
+ expected = artifacts(read_json(target / "inputs.json"))
+ for name, content in expected.items():
+ if (target / name).read_bytes() != content.encode():
+ raise ValueError(f"semantic replay mismatch: {name}")
+ return {"verified": True, "files": len(FILES), "independent_replay": True}
diff --git a/src/pitbridge/cli.py b/src/pitbridge/cli.py
new file mode 100644
index 0000000..65b7b42
--- /dev/null
+++ b/src/pitbridge/cli.py
@@ -0,0 +1,24 @@
+import argparse
+import json
+import sys
+
+from .bundle import read_json, verify_bundle, write_bundle
+from .demo import demo_inputs
+
+
+def main(argv=None):
+ parser = argparse.ArgumentParser(description="Build and independently replay availability-aware feature evidence")
+ sub = parser.add_subparsers(dest="command", required=True)
+ for command in ("demo", "build", "verify"):
+ p = sub.add_parser(command)
+ p.add_argument("--out", default="outputs")
+ if command == "build":
+ p.add_argument("--inputs", required=True, help="JSON observations, decisions and specs")
+ args = parser.parse_args(argv)
+ try:
+ result = verify_bundle(args.out) if args.command == "verify" else write_bundle(demo_inputs() if args.command == "demo" else read_json(args.inputs), args.out)
+ except (ValueError, OSError, KeyError) as exc:
+ print(f"pitbridge: {exc}", file=sys.stderr)
+ return 2
+ print(json.dumps(result, indent=2, ensure_ascii=False, allow_nan=False))
+ return 0
diff --git a/src/pitbridge/core.py b/src/pitbridge/core.py
new file mode 100644
index 0000000..f4c7041
--- /dev/null
+++ b/src/pitbridge/core.py
@@ -0,0 +1,209 @@
+"""Strict data contracts, a SQLite temporal join, and a separate replay oracle."""
+
+from dataclasses import asdict, dataclass
+from datetime import datetime, timezone
+import math
+import sqlite3
+from typing import Iterable
+
+
+def utc(value: str) -> str:
+ if not isinstance(value, str):
+ raise ValueError("timestamps must be timezone-aware ISO-8601 strings")
+ try:
+ parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
+ except ValueError as exc:
+ raise ValueError(f"invalid timestamp: {value!r}") from exc
+ if parsed.tzinfo is None or parsed.utcoffset() is None:
+ raise ValueError("naive timestamps are forbidden; include Z or an offset")
+ return parsed.astimezone(timezone.utc).isoformat(timespec="microseconds")
+
+
+def _identifier(value: str, name: str) -> None:
+ if not isinstance(value, str) or not value.strip() or value != value.strip():
+ raise ValueError(f"{name} must be a nonempty, trimmed string")
+
+
+@dataclass(frozen=True)
+class Observation:
+ record_id: str
+ entity_id: str
+ source: str
+ feature: str
+ event_at: str
+ published_at: str
+ ingested_at: str
+ revision: int
+ value: float | None
+ deleted: bool = False
+
+ def __post_init__(self):
+ for name in ("record_id", "entity_id", "source", "feature"):
+ _identifier(getattr(self, name), name)
+ for name in ("event_at", "published_at", "ingested_at"):
+ object.__setattr__(self, name, utc(getattr(self, name)))
+ if type(self.revision) is not int or not 1 <= self.revision <= 2**63-1:
+ raise ValueError("revision must be a positive int64 integer")
+ if type(self.deleted) is not bool:
+ raise ValueError("deleted must be a boolean")
+ if self.deleted:
+ if self.value is not None:
+ raise ValueError("a tombstone must have value=null")
+ elif isinstance(self.value, bool) or not isinstance(self.value, (int, float)) or not math.isfinite(self.value):
+ raise ValueError("active observations require a finite numeric value")
+ if self.published_at < self.event_at or self.ingested_at < self.event_at:
+ raise ValueError("publication and ingestion cannot precede the observation event")
+
+ @property
+ def available_at(self) -> str:
+ return max(self.published_at, self.ingested_at)
+
+
+@dataclass(frozen=True)
+class Decision:
+ decision_id: str
+ entity_id: str
+ decision_at: str
+
+ def __post_init__(self):
+ for name in ("decision_id", "entity_id"):
+ _identifier(getattr(self, name), name)
+ object.__setattr__(self, "decision_at", utc(self.decision_at))
+
+
+@dataclass(frozen=True)
+class FeatureSpec:
+ source: str
+ feature: str
+ max_age_days: float | None = None
+
+ def __post_init__(self):
+ for name in ("source", "feature"):
+ _identifier(getattr(self, name), name)
+ if self.max_age_days is not None and (isinstance(self.max_age_days, bool) or not isinstance(self.max_age_days, (int, float)) or not math.isfinite(self.max_age_days) or not 0 <= self.max_age_days <= 3652500):
+ raise ValueError("max_age_days must be within [0, 3652500], or null")
+
+
+@dataclass(frozen=True)
+class Snapshot:
+ decision_id: str
+ entity_id: str
+ decision_at: str
+ source: str
+ feature: str
+ value: float | None
+ status: str
+ record_id: str | None = None
+ revision: int | None = None
+ event_at: str | None = None
+ published_at: str | None = None
+ ingested_at: str | None = None
+ available_at: str | None = None
+
+
+def validate(observations: Iterable[Observation], decisions: Iterable[Decision], specs: Iterable[FeatureSpec]):
+ obs, dec, features = list(observations), list(decisions), list(specs)
+ for rows, cls in ((obs, Observation), (dec, Decision), (features, FeatureSpec)):
+ if any(not isinstance(row, cls) for row in rows):
+ raise ValueError(f"expected {cls.__name__} instances")
+ if not dec or not features:
+ raise ValueError("at least one decision and one feature spec are required")
+ for keys, name in (([r.record_id for r in obs], "record_id"), ([r.decision_id for r in dec], "decision_id"), ([(r.source, r.feature) for r in features], "feature spec"), ([(r.entity_id, r.source, r.feature, r.event_at, r.revision) for r in obs], "logical revision")):
+ if len(keys) != len(set(keys)):
+ raise ValueError(f"duplicate {name}")
+ allowed = {(r.source, r.feature) for r in features}
+ if any((r.source, r.feature) not in allowed for r in obs):
+ raise ValueError("every observation must have an explicit feature spec")
+ return sorted(obs, key=lambda r: r.record_id), sorted(dec, key=lambda r: r.decision_id), sorted(features, key=lambda r: (r.source, r.feature))
+
+
+def _instant(value: str) -> int:
+ delta = datetime.fromisoformat(value) - datetime(1970, 1, 1, tzinfo=timezone.utc)
+ return (delta.days * 86400 + delta.seconds) * 1_000_000 + delta.microseconds
+
+
+SQL = """
+WITH eligible AS (
+ SELECT d.decision_id, o.*,
+ ROW_NUMBER() OVER (
+ PARTITION BY d.decision_id, o.source, o.feature, o.event_at
+ ORDER BY o.revision DESC
+ ) AS version_rank
+ FROM decisions d JOIN observations o ON d.entity_id = o.entity_id
+ WHERE o.event_us <= d.decision_us AND o.available_us <= d.decision_us
+), current_versions AS (
+ SELECT * FROM eligible WHERE version_rank = 1
+), fresh AS (
+ SELECT c.*, ROW_NUMBER() OVER (
+ PARTITION BY c.decision_id, c.source, c.feature ORDER BY c.event_us DESC
+ ) AS event_rank
+ FROM current_versions c
+ JOIN decisions d USING (decision_id)
+ JOIN specs s USING (source, feature)
+ WHERE c.deleted = 0 AND
+ (s.max_age_us IS NULL OR d.decision_us - c.event_us <= s.max_age_us)
+)
+SELECT d.decision_id, d.entity_id, d.decision_at, s.source, s.feature,
+ f.value,
+ CASE WHEN f.record_id IS NOT NULL THEN 'selected'
+ WHEN NOT EXISTS (
+ SELECT 1 FROM observations o WHERE o.entity_id=d.entity_id
+ AND o.source=s.source AND o.feature=s.feature AND o.event_us<=d.decision_us
+ ) THEN 'no_history'
+ WHEN NOT EXISTS (
+ SELECT 1 FROM current_versions c WHERE c.decision_id=d.decision_id
+ AND c.source=s.source AND c.feature=s.feature
+ ) THEN 'not_available'
+ WHEN NOT EXISTS (
+ SELECT 1 FROM current_versions c WHERE c.decision_id=d.decision_id
+ AND c.source=s.source AND c.feature=s.feature AND c.deleted=0
+ ) THEN 'deleted'
+ ELSE 'stale' END AS status,
+ f.record_id, f.revision, f.event_at, f.published_at, f.ingested_at, f.available_at
+FROM decisions d CROSS JOIN specs s
+LEFT JOIN fresh f ON f.decision_id=d.decision_id AND f.source=s.source
+ AND f.feature=s.feature AND f.event_rank=1
+ORDER BY d.decision_id, s.source, s.feature
+"""
+
+
+def build_snapshot(observations, decisions, specs) -> list[Snapshot]:
+ """Use both event time and knowledge time, then resolve revisions and TTL."""
+ obs, dec, features = validate(observations, decisions, specs)
+ with sqlite3.connect(":memory:") as db:
+ db.executescript("""
+ CREATE TABLE observations(record_id TEXT,entity_id TEXT,source TEXT,feature TEXT,
+ event_at TEXT,published_at TEXT,ingested_at TEXT,revision INTEGER,value REAL,
+ deleted INTEGER,available_at TEXT,event_us INTEGER,available_us INTEGER);
+ CREATE TABLE decisions(decision_id TEXT,entity_id TEXT,decision_at TEXT,decision_us INTEGER);
+ CREATE TABLE specs(source TEXT,feature TEXT,max_age_us INTEGER);
+ CREATE INDEX temporal_lookup ON observations(entity_id,source,feature,event_us,available_us);
+ """)
+ db.executemany("INSERT INTO observations VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)", [tuple(asdict(r).values()) + (r.available_at, _instant(r.event_at), _instant(r.available_at)) for r in obs])
+ db.executemany("INSERT INTO decisions VALUES (?,?,?,?)", [(r.decision_id, r.entity_id, r.decision_at, _instant(r.decision_at)) for r in dec])
+ db.executemany("INSERT INTO specs VALUES (?,?,?)", [(r.source, r.feature, None if r.max_age_days is None else round(r.max_age_days * 86400 * 1e6)) for r in features])
+ return [Snapshot(*row) for row in db.execute(SQL)]
+
+
+def reference_snapshot(observations, decisions, specs) -> list[Snapshot]:
+ """Independent Python enumeration; intentionally does not call the SQL engine."""
+ obs, dec, features = validate(observations, decisions, specs)
+ results = []
+ for d in dec:
+ for s in features:
+ history = [o for o in obs if (o.entity_id, o.source, o.feature) == (d.entity_id, s.source, s.feature) and o.event_at <= d.decision_at]
+ known = [o for o in history if o.available_at <= d.decision_at]
+ versions = {}
+ for o in known:
+ if o.event_at not in versions or o.revision > versions[o.event_at].revision:
+ versions[o.event_at] = o
+ active = [o for o in versions.values() if not o.deleted]
+ age_limit = None if s.max_age_days is None else round(s.max_age_days * 86400 * 1e6)
+ fresh = [o for o in active if age_limit is None or _instant(d.decision_at) - _instant(o.event_at) <= age_limit]
+ status = "selected" if fresh else ("no_history" if not history else "not_available" if not known else "deleted" if not active else "stale")
+ row = Snapshot(d.decision_id, d.entity_id, d.decision_at, s.source, s.feature, None, status)
+ if fresh:
+ chosen = max(fresh, key=lambda o: o.event_at)
+ row = Snapshot(d.decision_id, d.entity_id, d.decision_at, s.source, s.feature, float(chosen.value), status, chosen.record_id, chosen.revision, chosen.event_at, chosen.published_at, chosen.ingested_at, chosen.available_at)
+ results.append(row)
+ return results
diff --git a/src/pitbridge/demo.py b/src/pitbridge/demo.py
new file mode 100644
index 0000000..851199d
--- /dev/null
+++ b/src/pitbridge/demo.py
@@ -0,0 +1,24 @@
+"""Compact, hand-auditable counterexamples across three financial sources."""
+
+
+def demo_inputs():
+ records = []
+ def add(rid, entity, source, feature, event, published, ingested, revision, value, deleted=False):
+ records.append(dict(record_id=rid, entity_id=entity, source=source, feature=feature, event_at=event+"T00:00:00Z", published_at=published+"T00:00:00Z", ingested_at=ingested+"T00:00:00Z", revision=revision, value=value, deleted=deleted))
+ add("tax-a-v1", "SME-A", "tax", "invoice_revenue", "2026-01-31", "2026-02-05", "2026-02-06", 1, 100000)
+ add("tax-a-v2", "SME-A", "tax", "invoice_revenue", "2026-01-31", "2026-03-10", "2026-03-11", 2, 180000)
+ add("tax-a-feb", "SME-A", "tax", "invoice_revenue", "2026-02-28", "2026-03-05", "2026-03-20", 1, 220000)
+ add("power-a", "SME-A", "utility", "kwh", "2026-02-01", "2026-02-02", "2026-02-03", 1, 1800)
+ add("cash-a", "SME-A", "bank", "net_inflow", "2026-02-10", "2026-02-10", "2026-02-11", 1, 50000)
+ add("cash-a-delete", "SME-A", "bank", "net_inflow", "2026-02-10", "2026-02-18", "2026-02-18", 2, None, True)
+ add("tax-b", "SME-B", "tax", "invoice_revenue", "2026-01-31", "2026-02-05", "2026-03-01", 1, 90000)
+ add("power-b-old", "SME-B", "utility", "kwh", "2025-11-01", "2025-11-02", "2025-11-02", 1, 700)
+ add("power-b", "SME-B", "utility", "kwh", "2026-02-01", "2026-02-03", "2026-03-02", 1, 1300)
+ add("cash-b-v2", "SME-B", "bank", "net_inflow", "2026-02-10", "2026-02-12", "2026-02-13", 2, 30000)
+ # A delayed delivery of revision 1 must not replace already-known revision 2.
+ add("cash-b-v1-delayed", "SME-B", "bank", "net_inflow", "2026-02-10", "2026-02-11", "2026-02-19", 1, 80000)
+ decisions = []
+ for entity in ("SME-A", "SME-B", "SME-C"):
+ for i, day in enumerate(("2026-02-04", "2026-02-15", "2026-02-20", "2026-03-15", "2026-03-25"), 1):
+ decisions.append(dict(decision_id=f"{entity}-{i}", entity_id=entity, decision_at=day+"T00:00:00Z"))
+ return {"data_kind": "synthetic", "observations": records, "decisions": decisions, "specs": [dict(source="tax", feature="invoice_revenue", max_age_days=90), dict(source="utility", feature="kwh", max_age_days=45), dict(source="bank", feature="net_inflow", max_age_days=60)]}
diff --git a/tests/test_bundle.py b/tests/test_bundle.py
new file mode 100644
index 0000000..86492bc
--- /dev/null
+++ b/tests/test_bundle.py
@@ -0,0 +1,82 @@
+from hashlib import sha256
+import json
+from pathlib import Path
+from tempfile import TemporaryDirectory
+import unittest
+
+from pitbridge.bundle import canonical, read_json, verify_bundle, write_bundle
+from pitbridge.cli import main
+from pitbridge.demo import demo_inputs
+
+
+class BundleTests(unittest.TestCase):
+ def test_demo_and_saved_inputs_replay(self):
+ with TemporaryDirectory() as tmp:
+ summary = write_bundle(demo_inputs(),tmp)
+ self.assertEqual(summary["future_knowledge_rows"],11)
+ self.assertEqual(summary["snapshot_rows"],45)
+ self.assertTrue(verify_bundle(tmp)["verified"])
+
+ def test_output_is_deterministic(self):
+ with TemporaryDirectory() as a, TemporaryDirectory() as b:
+ write_bundle(demo_inputs(),a); write_bundle(demo_inputs(),b)
+ self.assertEqual({p.name:p.read_bytes() for p in Path(a).iterdir()}, {p.name:p.read_bytes() for p in Path(b).iterdir()})
+
+ def test_changed_hash_fails_verification(self):
+ with TemporaryDirectory() as tmp:
+ write_bundle(demo_inputs(),tmp)
+ Path(tmp,"snapshots.csv").write_text("edited",encoding="utf-8")
+ with self.assertRaisesRegex(ValueError,"hash mismatch"):
+ verify_bundle(tmp)
+
+ def test_rehashed_false_metric_still_fails_semantic_replay(self):
+ with TemporaryDirectory() as tmp:
+ write_bundle(demo_inputs(),tmp)
+ path=Path(tmp,"summary.json")
+ data=read_json(path); data["future_knowledge_rows"]=0
+ path.write_text(canonical(data),encoding="utf-8")
+ manifest=read_json(Path(tmp,"manifest.json"))
+ manifest["files"][path.name]=sha256(path.read_bytes()).hexdigest()
+ Path(tmp,"manifest.json").write_text(canonical(manifest),encoding="utf-8")
+ with self.assertRaisesRegex(ValueError,"semantic replay mismatch"):
+ verify_bundle(tmp)
+
+ def test_unexpected_manifest_path_is_rejected(self):
+ with TemporaryDirectory() as tmp:
+ write_bundle(demo_inputs(),tmp)
+ manifest=read_json(Path(tmp,"manifest.json")); manifest["files"]["../outside"]="0"
+ Path(tmp,"manifest.json").write_text(canonical(manifest),encoding="utf-8")
+ with self.assertRaisesRegex(ValueError,"unexpected artifact"):
+ verify_bundle(tmp)
+
+ def test_nonfinite_and_duplicate_json_are_rejected(self):
+ with TemporaryDirectory() as tmp:
+ for text in ('{"a": NaN}', '{"a":1,"a":2}'):
+ Path(tmp,"bad.json").write_text(text,encoding="utf-8")
+ with self.assertRaises(ValueError):
+ read_json(Path(tmp,"bad.json"))
+
+ def test_html_escapes_untrusted_identifiers(self):
+ with TemporaryDirectory() as tmp:
+ data=demo_inputs(); data["decisions"][0]["decision_id"]=""
+ write_bundle(data,tmp)
+ html=Path(tmp,"report.html").read_text()
+ self.assertIn("<script>unsafe</script>",html)
+ self.assertNotIn("",html)
+
+ def test_user_supplied_report_is_not_labeled_synthetic(self):
+ with TemporaryDirectory() as tmp:
+ data=demo_inputs(); data.pop("data_kind")
+ write_bundle(data,tmp)
+ self.assertIn("USER-SUPPLIED DATA",Path(tmp,"report.html").read_text())
+
+ def test_cli_returns_nonzero_for_missing_artifacts(self):
+ with TemporaryDirectory() as tmp:
+ self.assertEqual(main(["verify","--out",tmp]),2)
+
+ def test_malformed_manifest_has_a_clear_validation_error(self):
+ with TemporaryDirectory() as tmp:
+ for payload in ([], {"files":None}, {"files":[]}):
+ Path(tmp,"manifest.json").write_text(canonical(payload),encoding="utf-8")
+ with self.assertRaisesRegex(ValueError,"unsupported manifest"):
+ verify_bundle(tmp)
diff --git a/tests/test_temporal.py b/tests/test_temporal.py
new file mode 100644
index 0000000..5c23ee7
--- /dev/null
+++ b/tests/test_temporal.py
@@ -0,0 +1,179 @@
+from dataclasses import asdict, replace
+from datetime import datetime, timedelta, timezone
+import random
+import unittest
+
+from pitbridge.core import Decision, FeatureSpec, Observation, build_snapshot, reference_snapshot
+from pitbridge.demo import demo_inputs
+from pitbridge.bundle import decode_inputs
+
+
+def observation(**changes):
+ fields = dict(record_id="r1", entity_id="A", source="tax", feature="revenue", event_at="2026-01-01T00:00:00Z", published_at="2026-01-02T00:00:00Z", ingested_at="2026-01-03T00:00:00Z", revision=1, value=10.0)
+ fields.update(changes)
+ return Observation(**fields)
+
+
+def snapshot(records, at="2026-01-05T00:00:00Z", age=None, entity="A"):
+ return build_snapshot(records, [Decision("D", entity, at)], [FeatureSpec("tax", "revenue", age)])[0]
+
+
+class TemporalTests(unittest.TestCase):
+ def test_historical_revision_cannot_change_old_decision(self):
+ revision = observation(record_id="r2", revision=2, value=90, published_at="2026-01-10T00:00:00Z", ingested_at="2026-01-10T00:00:00Z")
+ self.assertEqual(snapshot([observation(), revision]).value, 10)
+ self.assertEqual(snapshot([observation(), revision], "2026-01-10T00:00:00Z").value, 90)
+
+ def test_ingestion_delay_is_not_publication_time(self):
+ row = snapshot([observation(ingested_at="2026-01-08T00:00:00Z")])
+ self.assertEqual(row.status, "not_available")
+ self.assertIsNone(row.value)
+
+ def test_embargo_requires_publication_even_if_ingested(self):
+ self.assertEqual(snapshot([observation(published_at="2026-01-08T00:00:00Z")]).status, "not_available")
+
+ def test_equality_at_decision_is_inclusive(self):
+ row = snapshot([observation()], "2026-01-03T00:00:00Z")
+ self.assertEqual(row.record_id, "r1")
+
+ def test_microsecond_after_decision_is_excluded(self):
+ self.assertEqual(snapshot([observation()], "2026-01-02T23:59:59.999999Z").status, "not_available")
+
+ def test_future_event_is_never_selected(self):
+ self.assertEqual(snapshot([observation()], "2025-12-31T00:00:00Z").status, "no_history")
+
+ def test_max_age_boundary_includes_exact_cutoff(self):
+ self.assertEqual(snapshot([observation()], age=4).status, "selected")
+ self.assertEqual(snapshot([observation()], "2026-01-05T00:00:00.000001Z", age=4).status, "stale")
+
+ def test_zero_max_age_requires_same_instant(self):
+ row = observation(published_at="2026-01-01T00:00:00Z", ingested_at="2026-01-01T00:00:00Z")
+ self.assertEqual(snapshot([row], "2026-01-01T00:00:00Z", age=0).status, "selected")
+ self.assertEqual(snapshot([row], "2026-01-01T00:00:00.000001Z", age=0).status, "stale")
+
+ def test_known_tombstone_does_not_resurrect_old_revision(self):
+ tombstone = observation(record_id="del", revision=2, value=None, deleted=True)
+ self.assertEqual(snapshot([observation(), tombstone]).status, "deleted")
+
+ def test_future_tombstone_does_not_erase_history(self):
+ tombstone = observation(record_id="del", revision=2, value=None, deleted=True, ingested_at="2026-01-09T00:00:00Z")
+ self.assertEqual(snapshot([observation(), tombstone]).value, 10)
+
+ def test_tombstone_is_scoped_to_one_event(self):
+ recent = observation(record_id="recent", event_at="2026-01-02T00:00:00Z", value=30)
+ tombstone = replace(recent, record_id="del", revision=2, deleted=True, value=None)
+ self.assertEqual(snapshot([observation(), recent, tombstone]).record_id, "r1")
+
+ def test_late_lower_revision_does_not_replace_higher_revision(self):
+ old = observation(ingested_at="2026-01-04T00:00:00Z")
+ newer = observation(record_id="r2", revision=2, value=20)
+ self.assertEqual(snapshot([old, newer]).value, 20)
+
+ def test_event_order_precedes_revision_number(self):
+ old = observation(revision=99)
+ recent = observation(record_id="recent", event_at="2026-01-02T00:00:00Z", value=40)
+ self.assertEqual(snapshot([old, recent]).value, 40)
+
+ def test_timezones_are_compared_as_utc_instants(self):
+ row = observation(ingested_at="2026-01-03T08:00:00+08:00")
+ self.assertEqual(snapshot([row], "2026-01-02T19:00:00-05:00").status, "selected")
+
+ def test_entity_isolation_and_missing_rows_are_preserved(self):
+ self.assertEqual(snapshot([observation()], entity="B").status, "no_history")
+ rows = build_snapshot([], [Decision("D", "A", "2026-01-05T00:00:00Z")], [FeatureSpec("tax", "revenue"), FeatureSpec("bank", "cash")])
+ self.assertEqual(len(rows), 2)
+ self.assertTrue(all(r.value is None for r in rows))
+
+ def test_sources_with_same_feature_name_are_isolated(self):
+ obs = [observation(), observation(record_id="bank", source="bank", value=123)]
+ rows = build_snapshot(obs, [Decision("D", "A", "2026-01-05T00:00:00Z")], [FeatureSpec("tax", "revenue"), FeatureSpec("bank", "revenue")])
+ self.assertEqual({r.source:r.value for r in rows}, {"bank":123, "tax":10})
+
+ def test_shuffle_does_not_change_snapshot(self):
+ obs, dec, specs = decode_inputs(demo_inputs())
+ expected = build_snapshot(obs, dec, specs)
+ rng = random.Random(42)
+ for _ in range(10):
+ rng.shuffle(obs); rng.shuffle(dec); rng.shuffle(specs)
+ self.assertEqual(build_snapshot(obs,dec,specs), expected)
+
+ def test_appending_future_revisions_cannot_change_old_rows(self):
+ obs, dec, specs = decode_inputs(demo_inputs())
+ extra = replace(obs[0], record_id="future", revision=1000, value=None, deleted=True, published_at="2027-01-01T00:00:00Z", ingested_at="2027-01-01T00:00:00Z")
+ self.assertEqual(build_snapshot(obs,dec,specs), build_snapshot(obs+[extra],dec,specs))
+
+ def test_randomized_sql_matches_independent_enumeration(self):
+ rng = random.Random(711)
+ origin = datetime(2026,1,1,tzinfo=timezone.utc)
+ def stamp(day):
+ return (origin+timedelta(days=day)).isoformat()
+ for trial in range(40):
+ obs = []
+ for entity in ("A","B"):
+ for source in ("tax","bank"):
+ for event in range(4):
+ for revision in range(1,4):
+ deleted = rng.random() < .2
+ obs.append(Observation(f"{trial}-{entity}-{source}-{event}-{revision}",entity,source,"value",stamp(event*3),stamp(event*3+rng.randrange(8)),stamp(event*3+rng.randrange(12)),revision,None if deleted else rng.uniform(-100,100),deleted))
+ dec = [Decision(f"D{i}",rng.choice(("A","B","C")),stamp(rng.randrange(20))) for i in range(12)]
+ specs = [FeatureSpec("tax","value",rng.choice((None,0,5,20))),FeatureSpec("bank","value",5)]
+ self.assertEqual(build_snapshot(obs,dec,specs),reference_snapshot(obs,dec,specs))
+
+
+class ContractTests(unittest.TestCase):
+ def test_duplicate_ids_are_rejected(self):
+ with self.assertRaises(ValueError):
+ snapshot([observation(), observation()])
+
+ def test_duplicate_logical_revisions_are_rejected(self):
+ with self.assertRaises(ValueError):
+ snapshot([observation(), observation(record_id="other")])
+
+ def test_duplicate_decisions_and_specs_are_rejected(self):
+ d = Decision("D","A","2026-01-05T00:00:00Z")
+ s = FeatureSpec("tax","revenue")
+ for dec,spec in (([d,d],[s]),([d],[s,s])):
+ with self.subTest(decisions=len(dec),specs=len(spec)), self.assertRaises(ValueError):
+ build_snapshot([observation()],dec,spec)
+
+ def test_invalid_values_and_revisions_are_rejected(self):
+ for value in (float("nan"),float("inf"),None,True,"10"):
+ with self.subTest(value=value), self.assertRaises(ValueError):
+ observation(value=value)
+ for revision in (0,-1,1.5,True,2**63):
+ with self.subTest(revision=revision), self.assertRaises(ValueError):
+ observation(revision=revision)
+
+ def test_tombstone_requires_boolean_and_null_value(self):
+ for fields in ({"deleted":1},{"deleted":True,"value":10}):
+ with self.subTest(fields=fields), self.assertRaises(ValueError):
+ observation(**fields)
+
+ def test_naive_and_invalid_timestamps_are_rejected(self):
+ for stamp in ("2026-01-01", "bad", 0):
+ with self.subTest(stamp=stamp), self.assertRaises(ValueError):
+ observation(event_at=stamp)
+
+ def test_impossible_chronology_is_rejected(self):
+ for name in ("published_at","ingested_at"):
+ with self.subTest(field=name), self.assertRaises(ValueError):
+ observation(**{name:"2025-12-31T00:00:00Z"})
+
+ def test_negative_nonfinite_and_huge_ttl_are_rejected(self):
+ for age in (-1, True, float("inf"), 1e100):
+ with self.subTest(age=age), self.assertRaises(ValueError):
+ FeatureSpec("tax","revenue",age)
+
+ def test_unknown_features_are_not_silently_dropped(self):
+ with self.assertRaises(ValueError):
+ snapshot([observation(feature="unknown")])
+
+ def test_unknown_input_columns_are_rejected(self):
+ data = demo_inputs()
+ data["observations"][0]["future_label"] = 1
+ with self.assertRaises(ValueError):
+ decode_inputs(data)
+
+ def test_blank_ids_are_rejected(self):
+ with self.assertRaises(ValueError):
+ observation(entity_id=" ")