diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 85fc777..74feb77 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -23,8 +23,12 @@ jobs: - name: Counterexamples and contracts run: python -m unittest discover -s tests -v - name: Replay committed evidence - run: pitbridge verify --out demo + run: | + pitbridge verify --out demo + pitbridge verify-rolling --out demo/rolling - name: Generate and replay a fresh report run: | pitbridge demo --out outputs pitbridge verify --out outputs + pitbridge rolling-demo --out outputs/rolling + pitbridge verify-rolling --out outputs/rolling diff --git a/.github/workflows/pages.yml b/.github/workflows/pages.yml index 04dd580..819770a 100644 --- a/.github/workflows/pages.yml +++ b/.github/workflows/pages.yml @@ -20,11 +20,13 @@ jobs: run: | python -m pip install -e . pitbridge verify --out demo + pitbridge verify-rolling --out demo/rolling - name: Prepare the standalone report run: | mkdir _site - cp demo/* _site/ + cp -R demo/. _site/ cp demo/report.html _site/index.html + cp demo/rolling/report.html _site/rolling/index.html touch _site/.nojekyll - uses: actions/configure-pages@v5 - uses: actions/upload-pages-artifact@v4 diff --git a/README.md b/README.md index 9ed091c..c9ee4b8 100644 --- a/README.md +++ b/README.md @@ -63,8 +63,25 @@ The production path is a SQLite window-function join. A separate Python enumerat Verification checks hashes **and** reruns both algorithms, rejecting a fabricated summary even if its hash was updated. Hashes are not signatures and cannot establish that an external source is truthful. +## Availability-aware rolling features + +Rolling sums, observation means and counts resolve availability, revisions +and tombstones before aggregating an inclusive event-time window. Every result +retains all contributing records; a separate Python temporal enumerator checks +membership and values. + +```bash +pitbridge rolling-demo --out outputs/rolling +pitbridge verify-rolling --out outputs/rolling +pitbridge aggregate --inputs demo/rolling/inputs.json --out outputs/custom +``` + +[Online rolling case](https://dev-belly.github.io/PITBridge/rolling/) · +[Contract and worked example](docs/ROLLING.md) · +[Feature rows](demo/rolling/features.csv) · [Membership](demo/rolling/members.csv) + ## Scope -This is a reference implementation for scalar financial observations. It does not provide streaming ingestion, access control, label generation or aggregated rolling features. Version numbers and trustworthy availability timestamps must come from the upstream source contract. SQLite is intentionally inspectable; distributed-scale performance has not been benchmarked. +This is a reference implementation for scalar financial observations and explicit rolling aggregates. It does not provide streaming ingestion, access control or label generation. Version numbers and trustworthy availability timestamps must come from the upstream source contract. SQLite is intentionally inspectable; distributed-scale performance has not been benchmarked. Rolling means are observation means, and counts do not establish complete business activity. PITBridge complements [CreditVintage](https://github.com/dev-belly/CreditVintage)'s application-time model evaluation. The repositories are separate components; no integration is claimed. MIT license. diff --git a/README.zh-CN.md b/README.zh-CN.md index 1094cbc..96e0295 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -27,6 +27,14 @@ python -m unittest discover -s tests -v - **显式缺失:** 区分没有历史、尚未可用、已撤销和已过期,保留申请行。 - **逐条溯源:** 导出来源、记录 ID、版本、发生时间、公开时间与入库时间。 - **可重放证据:** 修改报告并更新哈希,仍会被语义重放发现。 +- **时点滚动统计:** 在选择当时已知版本后,计算指定窗口的加总、观察均值和计数,并保留每条贡献记录;SQL 与 Python 分别重放成员与数值。 + +[在线滚动特征案例](https://dev-belly.github.io/PITBridge/rolling/) · [窗口合同与手算案例](docs/ROLLING.md) + +```bash +pitbridge rolling-demo --out outputs/rolling +pitbridge verify-rolling --out outputs/rolling +``` 合成示例有 11 条源记录、15 次决策、45 次特征查询。错误的“只按发生时间关联”基线有 11 次使用未来信息;完整规则共改变 16 次选择。后者还包括撤销和过期影响,不能把两个数字混为一谈。 diff --git a/demo/manifest.json b/demo/manifest.json index 62e1258..7f728fe 100644 --- a/demo/manifest.json +++ b/demo/manifest.json @@ -3,7 +3,7 @@ "files": { "comparison.csv": "0eb458eceec7477682417da3bcfc24c1746021c52ad83b24c7a4a7e5fc4f734e", "inputs.json": "d55c44483897bd3ed838376b3e7a11f7d56ee41ef5cf8e6e22c99d060d701a98", - "report.html": "39b0dbcc2bd7c13447127990d4cece668fdbc1e15a180c969100dac8fa6d0f4c", + "report.html": "2d09d39b3e4034683a50c1fd566c27100aca20c891ec5b1124ad6713945cdf05", "snapshots.csv": "777f54dea6c8067cccb9a9e5ec375038a99dd2e44c9ab1df52cfae029a31e6f0", "summary.json": "fe3d30c39e8eb5244636e150f2da884e898456a1b5f52c2c42f0a65cf54e496b" }, diff --git a/demo/report.html b/demo/report.html index 18f28ac..24639df 100644 --- a/demo/report.html +++ b/demo/report.html @@ -5,6 +5,7 @@ *{box-sizing:border-box}body{margin:0;background:#0c1422;color:#e6edf7;font:16px/1.6 system-ui,sans-serif}main{max-width:1200px;margin:auto;padding:48px 24px}h1{font-size:clamp(32px,6vw,64px);line-height:1.1;margin:16px 0}h2{margin-top:44px}p{color:#a9b9cf}.tag{color:#64e5c4;letter-spacing:.14em;font-size:12px}.cards{display:grid;grid-template-columns:repeat(auto-fit,minmax(200px,1fr));gap:16px;margin:32px 0}.card{border:1px solid #293952;border-radius:14px;padding:22px;background:#111e31}.value{display:block;font-size:38px;font-weight:700}.scroll{overflow-x:auto;border:1px solid #293952;border-radius:12px}table{width:100%;border-collapse:collapse;font-size:13px}th,td{text-align:left;padding:12px;white-space:nowrap;border-bottom:1px solid #22314a}th{background:#17263b;color:#9fefd9}code{color:#9fefd9}a{color:#9fefd9}label{display:block;margin:16px 0;cursor:pointer}footer{margin-top:48px;color:#8ca1bc}.hidden{display:none}
PITBRIDGE / SYNTHETIC FINANCIAL DATA

What was known
when the decision was made?

+

Explore rolling-window features and their source membership

Availability-aware SQL snapshots with publication, ingestion, revision and record-level lineage.

15Decisions
45Feature lookups
11Unsafe future-data selections
16Selections changed

Independent replay: SQLite output matches the Python enumeration oracle. deleted: 3 · no_history: 17 · not_available: 4 · selected: 16 · stale: 5

diff --git a/demo/rolling/features.csv b/demo/rolling/features.csv new file mode 100644 index 0000000..42cff25 --- /dev/null +++ b/demo/rolling/features.csv @@ -0,0 +1,13 @@ +decision_id,entity_id,decision_at,name,source,feature,window_days,aggregation,value,observation_count,status +SME-A-10,SME-A,2026-02-10T18:00:00.000000+00:00,mean_inflow_7d,bank,net_inflow,7.0,mean,150.0,2,selected +SME-A-10,SME-A,2026-02-10T18:00:00.000000+00:00,net_inflow_30d,bank,net_inflow,30.0,sum,1500.0,4,selected +SME-A-10,SME-A,2026-02-10T18:00:00.000000+00:00,observations_30d,bank,net_inflow,30.0,count,4,4,selected +SME-A-11,SME-A,2026-02-11T18:00:00.000000+00:00,mean_inflow_7d,bank,net_inflow,7.0,mean,150.0,2,selected +SME-A-11,SME-A,2026-02-11T18:00:00.000000+00:00,net_inflow_30d,bank,net_inflow,30.0,sum,4300.0,3,selected +SME-A-11,SME-A,2026-02-11T18:00:00.000000+00:00,observations_30d,bank,net_inflow,30.0,count,3,3,selected +SME-A-12,SME-A,2026-02-12T18:00:00.000000+00:00,mean_inflow_7d,bank,net_inflow,7.0,mean,1250.0,2,selected +SME-A-12,SME-A,2026-02-12T18:00:00.000000+00:00,net_inflow_30d,bank,net_inflow,30.0,sum,6300.0,4,selected +SME-A-12,SME-A,2026-02-12T18:00:00.000000+00:00,observations_30d,bank,net_inflow,30.0,count,4,4,selected +SME-B-12,SME-B,2026-02-12T18:00:00.000000+00:00,mean_inflow_7d,bank,net_inflow,7.0,mean,,0,no_history +SME-B-12,SME-B,2026-02-12T18:00:00.000000+00:00,net_inflow_30d,bank,net_inflow,30.0,sum,,0,no_history +SME-B-12,SME-B,2026-02-12T18:00:00.000000+00:00,observations_30d,bank,net_inflow,30.0,count,,0,no_history diff --git a/demo/rolling/inputs.json b/demo/rolling/inputs.json new file mode 100644 index 0000000..098e72a --- /dev/null +++ b/demo/rolling/inputs.json @@ -0,0 +1,170 @@ +{ + "data_kind": "synthetic", + "decisions": [ + { + "decision_at": "2026-02-10T18:00:00.000000+00:00", + "decision_id": "SME-A-10", + "entity_id": "SME-A" + }, + { + "decision_at": "2026-02-11T18:00:00.000000+00:00", + "decision_id": "SME-A-11", + "entity_id": "SME-A" + }, + { + "decision_at": "2026-02-12T18:00:00.000000+00:00", + "decision_id": "SME-A-12", + "entity_id": "SME-A" + }, + { + "decision_at": "2026-02-12T18:00:00.000000+00:00", + "decision_id": "SME-B-12", + "entity_id": "SME-B" + } + ], + "observations": [ + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-01-11T18:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-01-11T18:00:00.000000+00:00", + "published_at": "2026-01-11T18:00:00.000000+00:00", + "record_id": "boundary", + "revision": 1, + "source": "bank", + "value": 200 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-07T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-07T00:00:00.000000+00:00", + "published_at": "2026-02-07T00:00:00.000000+00:00", + "record_id": "cancel-v1", + "revision": 1, + "source": "bank", + "value": 800 + }, + { + "deleted": true, + "entity_id": "SME-A", + "event_at": "2026-02-07T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-09T00:00:00.000000+00:00", + "published_at": "2026-02-09T00:00:00.000000+00:00", + "record_id": "cancel-v2", + "revision": 2, + "source": "bank", + "value": null + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-01T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-01T00:00:00.000000+00:00", + "published_at": "2026-02-01T00:00:00.000000+00:00", + "record_id": "cash-v1", + "revision": 1, + "source": "bank", + "value": 1000 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-01T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-11T00:00:00.000000+00:00", + "published_at": "2026-02-11T00:00:00.000000+00:00", + "record_id": "cash-v2", + "revision": 2, + "source": "bank", + "value": 4000 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-14T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-14T00:00:00.000000+00:00", + "published_at": "2026-02-14T00:00:00.000000+00:00", + "record_id": "future-event", + "revision": 1, + "source": "bank", + "value": 300 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-09T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-12T00:00:00.000000+00:00", + "published_at": "2026-02-12T00:00:00.000000+00:00", + "record_id": "late", + "revision": 1, + "source": "bank", + "value": 2000 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-01-11T17:59:59.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-01-11T17:59:59.000000+00:00", + "published_at": "2026-01-11T17:59:59.000000+00:00", + "record_id": "outside", + "revision": 1, + "source": "bank", + "value": 10000 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-08T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-09T00:00:00.000000+00:00", + "published_at": "2026-02-09T00:00:00.000000+00:00", + "record_id": "recent", + "revision": 1, + "source": "bank", + "value": 500 + }, + { + "deleted": false, + "entity_id": "SME-A", + "event_at": "2026-02-05T00:00:00.000000+00:00", + "feature": "net_inflow", + "ingested_at": "2026-02-05T00:00:00.000000+00:00", + "published_at": "2026-02-05T00:00:00.000000+00:00", + "record_id": "refund", + "revision": 1, + "source": "bank", + "value": -200 + } + ], + "rolling_specs": [ + { + "aggregation": "mean", + "feature": "net_inflow", + "name": "mean_inflow_7d", + "source": "bank", + "window_days": 7 + }, + { + "aggregation": "sum", + "feature": "net_inflow", + "name": "net_inflow_30d", + "source": "bank", + "window_days": 30 + }, + { + "aggregation": "count", + "feature": "net_inflow", + "name": "observations_30d", + "source": "bank", + "window_days": 30 + } + ] +} diff --git a/demo/rolling/manifest.json b/demo/rolling/manifest.json new file mode 100644 index 0000000..51ab126 --- /dev/null +++ b/demo/rolling/manifest.json @@ -0,0 +1,11 @@ +{ + "engine": "pitbridge/rolling/1", + "files": { + "features.csv": "62cbc8dfbd1cfa7567346341aa47b9067039186bd5a6a39d1e0f3f7771ad2695", + "inputs.json": "67480d63496734095c73eca104ddbf3843028be98238637b90facf8ee7859886", + "members.csv": "a5ebb522cc58139169d1655a02d6753a3f49f81f11517656c5152ebfa0251aa1", + "report.html": "8b7aa6e303fa6dbcceae3d630beb42d152673b2cf74aa1eca107099e7912316f", + "summary.json": "a2a15decb3ac74604c8a13dd1f1452695c35bfbab982d7c87ca234ba12fc8b1b" + }, + "schema_version": 1 +} diff --git a/demo/rolling/members.csv b/demo/rolling/members.csv new file mode 100644 index 0000000..ddebe6f --- /dev/null +++ b/demo/rolling/members.csv @@ -0,0 +1,29 @@ +decision_id,name,record_id,revision,event_at,published_at,ingested_at,available_at,value +SME-A-10,mean_inflow_7d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-10,mean_inflow_7d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-10,net_inflow_30d,boundary,1,2026-01-11T18:00:00.000000+00:00,2026-01-11T18:00:00.000000+00:00,2026-01-11T18:00:00.000000+00:00,2026-01-11T18:00:00.000000+00:00,200.0 +SME-A-10,net_inflow_30d,cash-v1,1,2026-02-01T00:00:00.000000+00:00,2026-02-01T00:00:00.000000+00:00,2026-02-01T00:00:00.000000+00:00,2026-02-01T00:00:00.000000+00:00,1000.0 +SME-A-10,net_inflow_30d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-10,net_inflow_30d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-10,observations_30d,boundary,1,2026-01-11T18:00:00.000000+00:00,2026-01-11T18:00:00.000000+00:00,2026-01-11T18:00:00.000000+00:00,2026-01-11T18:00:00.000000+00:00,200.0 +SME-A-10,observations_30d,cash-v1,1,2026-02-01T00:00:00.000000+00:00,2026-02-01T00:00:00.000000+00:00,2026-02-01T00:00:00.000000+00:00,2026-02-01T00:00:00.000000+00:00,1000.0 +SME-A-10,observations_30d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-10,observations_30d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-11,mean_inflow_7d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-11,mean_inflow_7d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-11,net_inflow_30d,cash-v2,2,2026-02-01T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,4000.0 +SME-A-11,net_inflow_30d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-11,net_inflow_30d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-11,observations_30d,cash-v2,2,2026-02-01T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,4000.0 +SME-A-11,observations_30d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-11,observations_30d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-12,mean_inflow_7d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-12,mean_inflow_7d,late,1,2026-02-09T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2000.0 +SME-A-12,net_inflow_30d,cash-v2,2,2026-02-01T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,4000.0 +SME-A-12,net_inflow_30d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-12,net_inflow_30d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-12,net_inflow_30d,late,1,2026-02-09T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2000.0 +SME-A-12,observations_30d,cash-v2,2,2026-02-01T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,2026-02-11T00:00:00.000000+00:00,4000.0 +SME-A-12,observations_30d,refund,1,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,2026-02-05T00:00:00.000000+00:00,-200.0 +SME-A-12,observations_30d,recent,1,2026-02-08T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,2026-02-09T00:00:00.000000+00:00,500.0 +SME-A-12,observations_30d,late,1,2026-02-09T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2026-02-12T00:00:00.000000+00:00,2000.0 diff --git a/demo/rolling/report.html b/demo/rolling/report.html new file mode 100644 index 0000000..e777812 --- /dev/null +++ b/demo/rolling/report.html @@ -0,0 +1,16 @@ + +PITBridge · Rolling financial features
PITBRIDGE / ROLLING FEATURES / SYNTHETIC DATA
+

What did the last 30 days
look like at decision time?

+

Sum, mean and count use the latest revision actually available for each event. Future corrections and known cancellations cannot inflate an earlier feature.

+
4Decisions
12Rolling lookups
28Contributing-record links
SQL = PythonIndependent temporal replay
+

Inspect the features

The window includes both boundaries. An empty or unavailable window keeps a missing value; its observation count is metadata, not an assumed business zero.

+ +
DecisionFeatureValueObservationsStatus
SME-A-10mean_inflow_7d150.02selected
SME-A-10net_inflow_30d1500.04selected
SME-A-10observations_30d44selected
SME-A-11mean_inflow_7d150.02selected
SME-A-11net_inflow_30d4300.03selected
SME-A-11observations_30d33selected
SME-A-12mean_inflow_7d1250.02selected
SME-A-12net_inflow_30d6300.04selected
SME-A-12observations_30d44selected
SME-B-12mean_inflow_7d—0no_history
SME-B-12net_inflow_30d—0no_history
SME-B-12observations_30d—0no_history
+

Trace every contribution

Each row identifies the selected source revision. One record can contribute to several declared windows. Counts refer to observations, not borrowers or days.

+
DecisionFeatureRecordRevisionEventAvailableValue
SME-A-10mean_inflow_7drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-10mean_inflow_7drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-10net_inflow_30dboundary12026-01-11T18:00:00.000000+00:002026-01-11T18:00:00.000000+00:00200.0
SME-A-10net_inflow_30dcash-v112026-02-01T00:00:00.000000+00:002026-02-01T00:00:00.000000+00:001000.0
SME-A-10net_inflow_30drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-10net_inflow_30drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-10observations_30dboundary12026-01-11T18:00:00.000000+00:002026-01-11T18:00:00.000000+00:00200.0
SME-A-10observations_30dcash-v112026-02-01T00:00:00.000000+00:002026-02-01T00:00:00.000000+00:001000.0
SME-A-10observations_30drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-10observations_30drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-11mean_inflow_7drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-11mean_inflow_7drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-11net_inflow_30dcash-v222026-02-01T00:00:00.000000+00:002026-02-11T00:00:00.000000+00:004000.0
SME-A-11net_inflow_30drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-11net_inflow_30drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-11observations_30dcash-v222026-02-01T00:00:00.000000+00:002026-02-11T00:00:00.000000+00:004000.0
SME-A-11observations_30drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-11observations_30drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-12mean_inflow_7drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-12mean_inflow_7dlate12026-02-09T00:00:00.000000+00:002026-02-12T00:00:00.000000+00:002000.0
SME-A-12net_inflow_30dcash-v222026-02-01T00:00:00.000000+00:002026-02-11T00:00:00.000000+00:004000.0
SME-A-12net_inflow_30drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-12net_inflow_30drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-12net_inflow_30dlate12026-02-09T00:00:00.000000+00:002026-02-12T00:00:00.000000+00:002000.0
SME-A-12observations_30dcash-v222026-02-01T00:00:00.000000+00:002026-02-11T00:00:00.000000+00:004000.0
SME-A-12observations_30drefund12026-02-05T00:00:00.000000+00:002026-02-05T00:00:00.000000+00:00-200.0
SME-A-12observations_30drecent12026-02-08T00:00:00.000000+00:002026-02-09T00:00:00.000000+00:00500.0
SME-A-12observations_30dlate12026-02-09T00:00:00.000000+00:002026-02-12T00:00:00.000000+00:002000.0
+

Reproduce the result

pitbridge verify-rolling --out demo/rolling checks hashes and reruns SQL and the Python enumerator for both features and membership. Binary64 values are not exact monetary accounting, and observed events do not establish complete source coverage.

+

Feature rows · Source membership · Inputs · Summary

+
+ diff --git a/demo/rolling/summary.json b/demo/rolling/summary.json new file mode 100644 index 0000000..a23eeb0 --- /dev/null +++ b/demo/rolling/summary.json @@ -0,0 +1,14 @@ +{ + "data_kind": "synthetic", + "decisions": 4, + "feature_rows": 12, + "independent_temporal_replay": true, + "member_rows": 28, + "observations": 10, + "schema_version": 1, + "status_counts": { + "no_history": 3, + "selected": 9 + }, + "window_boundaries": "inclusive" +} diff --git a/docs/INTERVIEW.md b/docs/INTERVIEW.md index 9d5b97a..8fe8837 100644 --- a/docs/INTERVIEW.md +++ b/docs/INTERVIEW.md @@ -28,4 +28,6 @@ ## 建议展示路径 +升级展示:打开 [滚动特征报告](https://dev-belly.github.io/PITBridge/rolling/),说明 30 天现金流从 1,500 到 4,300 再到 6,300 的变化如何同时受迟到、更正与窗口移动影响;在 `demo/rolling/members.csv` 逐条追溯。完整口径与追问见 [ROLLING.md](ROLLING.md)。 + 先打开 `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/ROLLING.md b/docs/ROLLING.md new file mode 100644 index 0000000..35385ca --- /dev/null +++ b/docs/ROLLING.md @@ -0,0 +1,94 @@ +# Decision-time rolling financial features + +[Online case](https://dev-belly.github.io/PITBridge/rolling/) · +[Feature rows](../demo/rolling/features.csv) · [Contributing records](../demo/rolling/members.csv) + +Rolling features answer: what was the sum, mean or observation count in a +recent event-time window, using only information available at the decision? + +## Run and replay + +```bash +pitbridge rolling-demo --out outputs/rolling +pitbridge verify-rolling --out outputs/rolling +pitbridge aggregate --inputs demo/rolling/inputs.json --out outputs/custom +``` + +Open `outputs/rolling/report.html`. Select a decision to inspect values and +every contributing record. Original scalar-snapshot commands remain supported. + +## Contract and selection order + +Inputs contain `observations`, `decisions`, `rolling_specs` and optional +`data_kind`. Observations and decisions retain their original strict contracts. + +```json +{"name":"net_inflow_30d","source":"bank","feature":"net_inflow", + "window_days":30,"aggregation":"sum"} +``` + +Output names are unique. Aggregation is `sum`, `mean` or `count`. Positive +window durations are rounded to microseconds. Both boundaries are inclusive: +`decision_at - window <= event_at <= decision_at`. + +1. Require event time and `max(published_at, ingested_at)` not later than the decision. +2. For each entity/source/feature/event-time key, retain the highest known revision. +3. Remove known tombstones after resolving revisions. +4. Aggregate active observations in the window, exporting their identities, + revisions and availability timestamps in `members.csv`. + +The SQLite temporal join uses a custom `math.fsum` aggregate. A separate +Python enumerator independently selects contributors. These paths share +compensated binary64 arithmetic, not their temporal selection logic. +This reduces cancellation error but is not exact decimal accounting; +non-finite outputs are rejected. Ordinary floating-point summation can +lose small terms when large values cancel; see the +[SQLite aggregate documentation](https://www.sqlite.org/lang_aggfunc.html). + +## Hand-auditable synthetic example + +Ten source revisions, four decisions and three specifications produce +12 feature rows and 28 contributing-record links. + +| Decision, 18:00 UTC | 30-day sum | Active observations | 7-day observation mean | +| --- | ---: | ---: | ---: | +| SME-A, February 10 | 1,500 | 4 | 150 | +| SME-A, February 11 | 4,300 | 3 | 150 | +| SME-A, February 12 | 6,300 | 4 | 1,250 | +| SME-B, February 12 | missing | 0 | missing | + +February 10 includes a 200-unit record exactly on the lower boundary, +excluding a record one second earlier, an unavailable 2,000-unit record +and a canceled 800-unit event. A 1,000-unit record is corrected to 4,000 +only on February 11, when the boundary record has already left the window. +The late 2,000-unit record first contributes on February 12. + +`no_history` means no events in this window; `not_available` means events +were not known; `deleted` means all known versions were canceled. These +keep `value=null` and `observation_count=0`. A selected sum of zero keeps +its real contributing records. Missing counts do not establish zero activity. + +## Evidence and boundaries + +The five-artifact bundle has a separate `pitbridge/rolling/1` manifest. +Verification checks hashes, then reruns SQL and Python for both values and +membership. Rehashed fabricated outputs are rejected. Hashes cannot +authenticate upstream records. + +The logical observation key is entity/source/feature/event timestamp. +Distinct transactions with the same key cannot be modeled independently. +Sources must supply comparable, nonoverlapping amounts for a meaningful sum. +Mean is an unweighted observation mean, not a complete daily/time-weighted +mean; count measures observations, not borrowers, invoices or days. +Complete source coverage, unit/currency conversion, streaming ingestion, +distributed throughput and CreditVintage integration remain outside scope. + +## 面试追问 + +**为什么不能先加总最近 30 天,再检查发布时间?** 未来修订已经改变了结果。必须先决定当时能看到哪个版本,再对合法成员聚合。 + +**为什么不能直接从原始表删掉撤销行?** 旧版本会重新被选中;这里先取最高已知修订,再处理撤销。 + +**没有记录为什么不是 0?** 没有历史、没收到数据和真实零经营额不同,必须保留缺失原因。 + +**怎么核验一笔加总?** 每个输出特征都有一组源修订成员,SQL 与 Python 从规范化输入分别重放,对数值和成员一起比较。 diff --git a/src/pitbridge/bundle.py b/src/pitbridge/bundle.py index e06c08d..f35d5be 100644 --- a/src/pitbridge/bundle.py +++ b/src/pitbridge/bundle.py @@ -98,6 +98,7 @@ def render_report(summary, comparison, snapshots): *{{box-sizing:border-box}}body{{margin:0;background:#0c1422;color:#e6edf7;font:16px/1.6 system-ui,sans-serif}}main{{max-width:1200px;margin:auto;padding:48px 24px}}h1{{font-size:clamp(32px,6vw,64px);line-height:1.1;margin:16px 0}}h2{{margin-top:44px}}p{{color:#a9b9cf}}.tag{{color:#64e5c4;letter-spacing:.14em;font-size:12px}}.cards{{display:grid;grid-template-columns:repeat(auto-fit,minmax(200px,1fr));gap:16px;margin:32px 0}}.card{{border:1px solid #293952;border-radius:14px;padding:22px;background:#111e31}}.value{{display:block;font-size:38px;font-weight:700}}.scroll{{overflow-x:auto;border:1px solid #293952;border-radius:12px}}table{{width:100%;border-collapse:collapse;font-size:13px}}th,td{{text-align:left;padding:12px;white-space:nowrap;border-bottom:1px solid #22314a}}th{{background:#17263b;color:#9fefd9}}code{{color:#9fefd9}}a{{color:#9fefd9}}label{{display:block;margin:16px 0;cursor:pointer}}footer{{margin-top:48px;color:#8ca1bc}}.hidden{{display:none}}
PITBRIDGE / {kind}

What was known
when the decision was made?

+

Explore rolling-window features and their source membership

Availability-aware SQL snapshots with publication, ingestion, revision and record-level lineage.

{summary['decisions']}Decisions
{summary['snapshot_rows']}Feature lookups
{summary['future_knowledge_rows']}Unsafe future-data selections
{summary['different_selections']}Selections changed

Independent replay: SQLite output matches the Python enumeration oracle. {counts}

diff --git a/src/pitbridge/cli.py b/src/pitbridge/cli.py index 65b7b42..092c822 100644 --- a/src/pitbridge/cli.py +++ b/src/pitbridge/cli.py @@ -4,19 +4,26 @@ from .bundle import read_json, verify_bundle, write_bundle from .demo import demo_inputs +from .rolling_bundle import rolling_demo_inputs, verify_rolling_bundle, write_rolling_bundle 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"): + for command in ("demo", "build", "verify", "rolling-demo", "aggregate", "verify-rolling"): p = sub.add_parser(command) - p.add_argument("--out", default="outputs") - if command == "build": + p.add_argument("--out", default="outputs/rolling" if command in ("rolling-demo", "aggregate", "verify-rolling") else "outputs") + if command in ("build", "aggregate"): 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) + if args.command == "verify-rolling": + result = verify_rolling_bundle(args.out) + elif args.command in ("rolling-demo", "aggregate"): + data = rolling_demo_inputs() if args.command == "rolling-demo" else read_json(args.inputs) + result = write_rolling_bundle(data, args.out) + else: + 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 diff --git a/src/pitbridge/rolling.py b/src/pitbridge/rolling.py new file mode 100644 index 0000000..67afb8a --- /dev/null +++ b/src/pitbridge/rolling.py @@ -0,0 +1,203 @@ +"""Availability-aware rolling aggregates with every contributing revision retained.""" + +from dataclasses import asdict, dataclass +import math +import sqlite3 + +from .core import Decision, FeatureSpec, Observation, _identifier, _instant, validate + + +@dataclass(frozen=True) +class RollingSpec: + name: str + source: str + feature: str + window_days: float + aggregation: str = "sum" + + def __post_init__(self): + for field in ("name", "source", "feature"): + _identifier(getattr(self, field), field) + if (isinstance(self.window_days, bool) or not isinstance(self.window_days, (int, float)) + or not math.isfinite(self.window_days) or not 0 < self.window_days <= 3652500 + or self.window_us < 1): + raise ValueError("window_days must be positive and resolve to at least one microsecond") + if self.aggregation not in ("sum", "mean", "count"): + raise ValueError("aggregation must be sum, mean or count") + + @property + def window_us(self): + return round(self.window_days * 86400 * 1_000_000) + + +@dataclass(frozen=True) +class RollingFeature: + decision_id: str + entity_id: str + decision_at: str + name: str + source: str + feature: str + window_days: float + aggregation: str + value: float | int | None + observation_count: int + status: str + + +@dataclass(frozen=True) +class RollingMember: + decision_id: str + name: str + record_id: str + revision: int + event_at: str + published_at: str + ingested_at: str + available_at: str + value: float + + +def validate_rolling(observations, decisions, specs): + specs = list(specs) + if not specs or any(not isinstance(spec, RollingSpec) for spec in specs): + raise ValueError("at least one RollingSpec is required") + if len({spec.name for spec in specs}) != len(specs): + raise ValueError("rolling output names must be unique") + sources = sorted({(spec.source, spec.feature) for spec in specs}) + observations, decisions, _ = validate( + observations, decisions, [FeatureSpec(*source) for source in sources] + ) + return observations, decisions, sorted(specs, key=lambda spec: spec.name) + + +class _PreciseSum: + """Compensated binary64 summation; this is not exact decimal accounting.""" + def __init__(self): + self.values = [] + + def step(self, value): + if value is not None: + self.values.append(float(value)) + + def finalize(self): + try: + return math.fsum(self.values) if self.values else None + except OverflowError: + return math.inf # Rejected by the finite-output check below. + + +CTE = """ +WITH eligible AS ( + SELECT d.decision_id, o.*, + ROW_NUMBER() OVER ( + PARTITION BY d.decision_id, o.source, o.feature, o.event_us + 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 +), active AS ( + SELECT e.*, s.name + FROM eligible e JOIN decisions d USING (decision_id) + JOIN specs s USING (source, feature) + WHERE e.version_rank = 1 AND e.deleted = 0 + AND d.decision_us - e.event_us <= s.window_us +) +""" + + +def build_rolling(observations, decisions, specs): + """Resolve known revisions before aggregating the inclusive event-time window.""" + observations, decisions, specs = validate_rolling(observations, decisions, specs) + with sqlite3.connect(":memory:") as db: + db.create_aggregate("precise_sum", 1, _PreciseSum) + 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(name TEXT,source TEXT,feature TEXT,window_days REAL, + aggregation TEXT,window_us INTEGER); + CREATE INDEX temporal_lookup ON observations(entity_id,source,feature,event_us,available_us); + """) + db.executemany("INSERT INTO observations VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)", [ + tuple(asdict(row).values()) + (row.available_at, _instant(row.event_at), _instant(row.available_at)) + for row in observations + ]) + db.executemany("INSERT INTO decisions VALUES (?,?,?,?)", [ + (row.decision_id, row.entity_id, row.decision_at, _instant(row.decision_at)) + for row in decisions + ]) + db.executemany("INSERT INTO specs VALUES (?,?,?,?,?,?)", [ + tuple(asdict(spec).values()) + (spec.window_us,) for spec in specs + ]) + query = CTE + """ + SELECT d.decision_id,d.entity_id,d.decision_at,s.name,s.source,s.feature, + s.window_days,s.aggregation, + CASE WHEN count(a.record_id)=0 THEN NULL + WHEN s.aggregation='count' THEN count(a.record_id) + WHEN s.aggregation='sum' THEN precise_sum(CASE WHEN s.aggregation='count' THEN NULL ELSE a.value END) + ELSE precise_sum(CASE WHEN s.aggregation='count' THEN NULL ELSE a.value END)/count(a.record_id) END, + count(a.record_id), + CASE WHEN count(a.record_id)>0 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 AND d.decision_us-o.event_us<=s.window_us) + THEN 'no_history' + WHEN NOT EXISTS (SELECT 1 FROM eligible e WHERE e.decision_id=d.decision_id + AND e.source=s.source AND e.feature=s.feature AND e.version_rank=1 + AND d.decision_us-e.event_us<=s.window_us) THEN 'not_available' + ELSE 'deleted' END + FROM decisions d CROSS JOIN specs s + LEFT JOIN active a ON a.decision_id=d.decision_id AND a.name=s.name + GROUP BY d.decision_id,s.name ORDER BY d.decision_id,s.name + """ + try: + features = [RollingFeature(*row) for row in db.execute(query)] + members = [RollingMember(*row) for row in db.execute(CTE + """ + SELECT decision_id,name,record_id,revision,event_at,published_at, + ingested_at,available_at,value FROM active + ORDER BY decision_id,name,event_us,record_id + """)] + except sqlite3.DataError as error: + raise ValueError("rolling aggregate is outside the finite binary64 range") from error + if any(row.value is not None and not math.isfinite(row.value) for row in features): + raise ValueError("rolling aggregate is outside the finite binary64 range") + return features, members + + +def reference_rolling(observations, decisions, specs): + """Independent temporal enumeration, sharing only the declared arithmetic primitive.""" + observations, decisions, specs = validate_rolling(observations, decisions, specs) + features, members = [], [] + for decision in decisions: + for spec in specs: + history = [row for row in observations + if (row.entity_id, row.source, row.feature) == (decision.entity_id, spec.source, spec.feature) + and 0 <= _instant(decision.decision_at) - _instant(row.event_at) <= spec.window_us] + known = [row for row in history if row.available_at <= decision.decision_at] + versions = {} + for row in known: + if row.event_at not in versions or versions[row.event_at].revision < row.revision: + versions[row.event_at] = row + active = sorted((row for row in versions.values() if not row.deleted), + key=lambda row: (row.event_at, row.record_id)) + status = "selected" if active else "no_history" if not history else "not_available" if not known else "deleted" + value = None + if active: + if spec.aggregation == "count": + value = len(active) + else: + try: + value = math.fsum(float(row.value) for row in active) + except OverflowError as error: + raise ValueError("rolling aggregate is outside the finite binary64 range") from error + if spec.aggregation == "mean": + value /= len(active) + features.append(RollingFeature(decision.decision_id, decision.entity_id, + decision.decision_at, spec.name, spec.source, spec.feature, spec.window_days, + spec.aggregation, value, len(active), status)) + members.extend(RollingMember(decision.decision_id, spec.name, row.record_id, + row.revision, row.event_at, row.published_at, row.ingested_at, + row.available_at, float(row.value)) for row in active) + return features, members diff --git a/src/pitbridge/rolling_bundle.py b/src/pitbridge/rolling_bundle.py new file mode 100644 index 0000000..59d7ae3 --- /dev/null +++ b/src/pitbridge/rolling_bundle.py @@ -0,0 +1,130 @@ +"""A separately versioned rolling-feature bundle; scalar bundles remain unchanged.""" + +from collections import Counter +from dataclasses import asdict, fields +from hashlib import sha256 +from html import escape +from pathlib import Path + +from .bundle import canonical, csv_text, read_json +from .core import Decision, Observation +from .rolling import RollingFeature, RollingMember, RollingSpec, build_rolling, reference_rolling, validate_rolling + +FILES = ("inputs.json", "features.csv", "members.csv", "summary.json", "report.html") +ENGINE = "pitbridge/rolling/1" + + +def decode_rolling(data): + required = {"observations", "decisions", "rolling_specs"} + if not isinstance(data, dict) or not required <= data.keys() or data.keys() - required - {"data_kind"}: + raise ValueError("rolling inputs require observations, decisions, rolling_specs; only data_kind is optional") + if data.get("data_kind", "user_supplied") not in ("synthetic", "user_supplied"): + raise ValueError("invalid data_kind") + if any(not isinstance(data[key], list) for key in required): + raise ValueError("rolling input sections must be arrays") + try: + return validate_rolling([Observation(**row) for row in data["observations"]], + [Decision(**row) for row in data["decisions"]], + [RollingSpec(**row) for row in data["rolling_specs"]]) + except TypeError as error: + raise ValueError(f"rolling input schema mismatch: {error}") from error + + +def rolling_demo_inputs(): + def observation(record_id, event_at, value, revision=1, available_at=None, deleted=False): + available_at = available_at or event_at + return dict(record_id=record_id, entity_id="SME-A", source="bank", feature="net_inflow", + event_at=event_at, published_at=available_at, ingested_at=available_at, + revision=revision, value=value, deleted=deleted) + return {"data_kind": "synthetic", "observations": [ + observation("boundary", "2026-01-11T18:00:00Z", 200), + observation("outside", "2026-01-11T17:59:59Z", 10000), + observation("cash-v1", "2026-02-01T00:00:00Z", 1000), + observation("cash-v2", "2026-02-01T00:00:00Z", 4000, 2, "2026-02-11T00:00:00Z"), + observation("refund", "2026-02-05T00:00:00Z", -200), + observation("recent", "2026-02-08T00:00:00Z", 500, available_at="2026-02-09T00:00:00Z"), + observation("cancel-v1", "2026-02-07T00:00:00Z", 800), + observation("cancel-v2", "2026-02-07T00:00:00Z", None, 2, "2026-02-09T00:00:00Z", True), + observation("late", "2026-02-09T00:00:00Z", 2000, available_at="2026-02-12T00:00:00Z"), + observation("future-event", "2026-02-14T00:00:00Z", 300), + ], "decisions": [ + dict(decision_id=f"SME-A-{day}", entity_id="SME-A", decision_at=f"2026-02-{day}T18:00:00Z") + for day in (10, 11, 12) + ] + [dict(decision_id="SME-B-12", entity_id="SME-B", decision_at="2026-02-12T18:00:00Z")], + "rolling_specs": [ + dict(name="net_inflow_30d", source="bank", feature="net_inflow", window_days=30, aggregation="sum"), + dict(name="mean_inflow_7d", source="bank", feature="net_inflow", window_days=7, aggregation="mean"), + dict(name="observations_30d", source="bank", feature="net_inflow", window_days=30, aggregation="count"), + ]} + + +def render_rolling(summary, features, members): + label = "SYNTHETIC DATA" if summary["data_kind"] == "synthetic" else "USER-SUPPLIED DATA" + def table_rows(rows, columns): + return "".join(f'' + + "".join(f"{escape(str(getattr(row, key) if getattr(row, key) is not None else '—'))}" + for key in columns) + "" for row in rows) + feature_rows = table_rows(features, ("decision_id", "name", "value", "observation_count", "status")) + member_rows = table_rows(members, ("decision_id", "name", "record_id", "revision", "event_at", "available_at", "value")) + options = "".join(f'' for name in sorted({row.decision_id for row in features})) + return f""" +PITBridge · Rolling financial features
PITBRIDGE / ROLLING FEATURES / {label}
+

What did the last 30 days
look like at decision time?

+

Sum, mean and count use the latest revision actually available for each event. Future corrections and known cancellations cannot inflate an earlier feature.

+
{summary['decisions']}Decisions
{summary['feature_rows']}Rolling lookups
{summary['member_rows']}Contributing-record links
SQL = PythonIndependent temporal replay
+

Inspect the features

The window includes both boundaries. An empty or unavailable window keeps a missing value; its observation count is metadata, not an assumed business zero.

+ +
{feature_rows}
DecisionFeatureValueObservationsStatus
+

Trace every contribution

Each row identifies the selected source revision. One record can contribute to several declared windows. Counts refer to observations, not borrowers or days.

+
{member_rows}
DecisionFeatureRecordRevisionEventAvailableValue
+

Reproduce the result

pitbridge verify-rolling --out demo/rolling checks hashes and reruns SQL and the Python enumerator for both features and membership. Binary64 values are not exact monetary accounting, and observed events do not establish complete source coverage.

+

Feature rows · Source membership · Inputs · Summary

+
+\n""" + + +def rolling_artifacts(data): + observations, decisions, specs = decode_rolling(data) + features, members = build_rolling(observations, decisions, specs) + if (features, members) != reference_rolling(observations, decisions, specs): + raise ValueError("rolling SQL/reference mismatch") + normalized = {"data_kind": data.get("data_kind", "user_supplied"), + "observations": [asdict(row) for row in observations], + "decisions": [asdict(row) for row in decisions], "rolling_specs": [asdict(spec) for spec in specs]} + summary = {"schema_version": 1, "data_kind": normalized["data_kind"], "observations": len(observations), + "decisions": len(decisions), "feature_rows": len(features), "member_rows": len(members), + "status_counts": dict(sorted(Counter(row.status for row in features).items())), + "window_boundaries": "inclusive", "independent_temporal_replay": True} + return {"inputs.json": canonical(normalized), "summary.json": canonical(summary), + "features.csv": csv_text([asdict(row) for row in features], [field.name for field in fields(RollingFeature)]), + "members.csv": csv_text([asdict(row) for row in members], [field.name for field in fields(RollingMember)]), + "report.html": render_rolling(summary, features, members)} + + +def write_rolling_bundle(data, out): + artifacts = rolling_artifacts(data) + target = Path(out) + target.mkdir(parents=True, exist_ok=True) + for name, content in artifacts.items(): + (target/name).write_text(content, encoding="utf-8", newline="") + manifest = {"schema_version": 1, "engine": ENGINE, + "files": {name: sha256(content.encode()).hexdigest() for name, content in sorted(artifacts.items())}} + (target/"manifest.json").write_text(canonical(manifest), encoding="utf-8") + return read_json(target/"summary.json") + + +def verify_rolling_bundle(out): + target = Path(out) + manifest = read_json(target/"manifest.json") + if (not isinstance(manifest, dict) or manifest.get("schema_version") != 1 or manifest.get("engine") != ENGINE + or not isinstance(manifest.get("files"), dict) or set(manifest["files"]) != set(FILES)): + raise ValueError("unsupported rolling 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}") + for name, content in rolling_artifacts(read_json(target/"inputs.json")).items(): + if (target/name).read_bytes() != content.encode(): + raise ValueError(f"rolling semantic replay mismatch: {name}") + return {"verified": True, "files": len(FILES), "independent_temporal_replay": True} diff --git a/tests/test_rolling.py b/tests/test_rolling.py new file mode 100644 index 0000000..761ad7d --- /dev/null +++ b/tests/test_rolling.py @@ -0,0 +1,237 @@ +from dataclasses import replace +from datetime import datetime, timedelta, timezone +from hashlib import sha256 +import contextlib +import io +import json +from pathlib import Path +import random +import tempfile +import unittest + +from pitbridge.bundle import canonical, read_json +from pitbridge.cli import main +from pitbridge.core import Decision, Observation +from pitbridge.rolling import RollingSpec, build_rolling, reference_rolling +from pitbridge.rolling_bundle import decode_rolling, rolling_artifacts, rolling_demo_inputs, verify_rolling_bundle, write_rolling_bundle + + +def observation(record_id, event, value, revision=1, published=None, ingested=None, deleted=False): + return Observation(record_id, "A", "bank", "cash", event, published or event, + ingested or published or event, revision, value, deleted) + + +def spec(name="cash_30d", days=30, aggregation="sum"): + return RollingSpec(name, "bank", "cash", days, aggregation) + + +class RollingTests(unittest.TestCase): + def test_hand_auditable_demo_sums_and_members(self): + data = decode_rolling(rolling_demo_inputs()) + features, members = build_rolling(*data) + sums = {row.decision_id: row.value for row in features if row.name == "net_inflow_30d"} + self.assertEqual(sums, {"SME-A-10": 1500, "SME-A-11": 4300, "SME-A-12": 6300, "SME-B-12": None}) + first = {row.record_id for row in members if row.decision_id == "SME-A-10" and row.name == "net_inflow_30d"} + self.assertEqual(first, {"boundary", "cash-v1", "refund", "recent"}) + self.assertEqual((features, members), reference_rolling(*data)) + + def test_both_window_boundaries_are_inclusive_to_microseconds(self): + rows = [observation("lower", "2026-01-01T00:00:00Z", 10), + observation("too-old", "2025-12-31T23:59:59.999999Z", 100), + observation("upper", "2026-01-02T00:00:00Z", 2), + observation("future", "2026-01-02T00:00:00.000001Z", 1000)] + result, members = build_rolling(rows, [Decision("D", "A", "2026-01-02T00:00:00Z")], [spec(days=1)]) + self.assertEqual(result[0].value, 12) + self.assertEqual({row.record_id for row in members}, {"lower", "upper"}) + + def test_later_ingestion_controls_availability(self): + row = observation("late", "2026-01-01T00:00:00Z", 10, + published="2026-01-02T00:00:00Z", ingested="2026-01-04T00:00:00Z") + result, members = build_rolling([row], [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec()]) + self.assertEqual(result[0].status, "not_available") + self.assertIsNone(result[0].value) + self.assertEqual(members, []) + + def test_later_publication_also_controls_availability(self): + row = observation("late", "2026-01-01T00:00:00Z", 10, + published="2026-01-04T00:00:00Z", ingested="2026-01-02T00:00:00Z") + result, _ = build_rolling([row], [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec()]) + self.assertEqual(result[0].status, "not_available") + + def test_highest_known_revision_wins_over_last_arrival(self): + rows = [observation("v2", "2026-01-01T00:00:00Z", 20, 2, "2026-01-02T00:00:00Z"), + observation("v1-late", "2026-01-01T00:00:00Z", 10, 1, "2026-01-03T00:00:00Z")] + result, members = build_rolling(rows, [Decision("D", "A", "2026-01-04T00:00:00Z")], [spec()]) + self.assertEqual(result[0].value, 20) + self.assertEqual([row.record_id for row in members], ["v2"]) + + def test_tombstone_removes_one_event_without_resurrecting_its_revision(self): + rows = [observation("older-event", "2026-01-01T00:00:00Z", 10), + observation("v1", "2026-01-02T00:00:00Z", 20), + observation("v2", "2026-01-02T00:00:00Z", None, 2, "2026-01-03T00:00:00Z", deleted=True)] + result, members = build_rolling(rows, [Decision("D", "A", "2026-01-04T00:00:00Z")], [spec()]) + self.assertEqual(result[0].value, 10) + self.assertEqual([row.record_id for row in members], ["older-event"]) + + def test_future_historical_correction_does_not_change_earlier_features(self): + rows = [observation("v1", "2026-01-01T00:00:00Z", 10)] + decision = [Decision("D", "A", "2026-01-03T00:00:00Z")] + before = build_rolling(rows, decision, [spec()]) + rows.append(observation("v2", "2026-01-01T00:00:00Z", 200, 2, "2026-02-01T00:00:00Z")) + self.assertEqual(before, build_rolling(rows, decision, [spec()])) + + def test_real_zero_and_negative_values_remain_selected(self): + rows = [observation("income", "2026-01-01T00:00:00Z", 100), + observation("refund", "2026-01-02T00:00:00Z", -100)] + result, _ = build_rolling(rows, [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec()]) + self.assertEqual((result[0].value, result[0].status, result[0].observation_count), (0, "selected", 2)) + + def test_empty_count_is_missing_not_an_imputed_business_zero(self): + result, members = build_rolling([], [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec(aggregation="count")]) + self.assertEqual((result[0].value, result[0].observation_count, result[0].status), (None, 0, "no_history")) + self.assertEqual(members, []) + + def test_fully_deleted_window_has_an_explicit_reason(self): + rows = [observation("v1", "2026-01-01T00:00:00Z", 10), + observation("v2", "2026-01-01T00:00:00Z", None, 2, "2026-01-02T00:00:00Z", deleted=True)] + result, _ = build_rolling(rows, [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec()]) + self.assertEqual((result[0].value, result[0].status), (None, "deleted")) + + def test_multiple_windows_have_distinct_membership_and_means(self): + rows = [observation("old", "2026-01-01T00:00:00Z", 100), + observation("new", "2026-01-09T00:00:00Z", 10)] + specs = [spec("sum", 30), spec("mean", 30, "mean"), spec("count", 1, "count")] + features, members = build_rolling(rows, [Decision("D", "A", "2026-01-10T00:00:00Z")], specs) + self.assertEqual({row.name: row.value for row in features}, {"sum": 110, "mean": 55, "count": 1}) + self.assertEqual([row.record_id for row in members if row.name == "count"], ["new"]) + + def test_invalid_window_contracts_are_rejected(self): + for days in (0, -1, True, float("nan"), float("inf"), 1e-20, 3652501): + with self.subTest(days=days), self.assertRaises(ValueError): + spec(days=days) + with self.assertRaises(ValueError): + spec(aggregation="median") + + def test_duplicate_output_names_are_rejected(self): + with self.assertRaisesRegex(ValueError, "unique"): + build_rolling([], [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec(), spec(days=7)]) + + def test_duplicate_source_revision_is_rejected(self): + row = observation("v1", "2026-01-01T00:00:00Z", 10) + with self.assertRaisesRegex(ValueError, "logical revision"): + build_rolling([row, replace(row, record_id="duplicate")], [Decision("D", "A", "2026-01-03T00:00:00Z")], [spec()]) + + def test_compensated_sum_survives_large_cancellation(self): + rows = [observation(str(i), f"2026-01-0{i+1}T00:00:00Z", value) + for i, value in enumerate((1e16, 1, -1e16))] + data = (rows, [Decision("D", "A", "2026-01-04T00:00:00Z")], [spec()]) + self.assertEqual(build_rolling(*data)[0][0].value, 1) + self.assertEqual(build_rolling(*data), reference_rolling(*data)) + + def test_sum_overflow_is_rejected_but_count_is_valid(self): + rows = [observation(str(i), f"2026-01-0{i+1}T00:00:00Z", 1e308) for i in range(2)] + decisions = [Decision("D", "A", "2026-01-04T00:00:00Z")] + with self.assertRaisesRegex(ValueError, "binary64"): + build_rolling(rows, decisions, [spec()]) + with self.assertRaisesRegex(ValueError, "binary64"): + reference_rolling(rows, decisions, [spec()]) + self.assertEqual(build_rolling(rows, decisions, [spec(aggregation="count")])[0][0].value, 2) + + def test_randomized_revision_histories_match_independent_temporal_oracle(self): + rng = random.Random(20261002) + anchor = datetime(2026, 2, 1, tzinfo=timezone.utc) + for trial in range(60): + rows = [] + for entity in ("A", "B"): + for event in range(8): + time = anchor + timedelta(days=rng.randint(-35, 10), seconds=event) + for revision in range(1, rng.randint(1, 3)+1): + deleted = rng.random() < 0.2 + rows.append(Observation(f"{entity}-{event}-{revision}", entity, "bank", "cash", + time.isoformat(), (time+timedelta(days=rng.randint(0, 12))).isoformat(), + (time+timedelta(days=rng.randint(0, 15))).isoformat(), revision, + None if deleted else rng.randint(-100, 200)/2, deleted)) + decisions = [Decision(f"{entity}-{day}", entity, (anchor+timedelta(days=day)).isoformat()) + for entity in ("A", "B") for day in (0, 5, 10)] + specs = [spec("sum", 30), spec("mean", 7, "mean"), spec("count", 1, "count")] + with self.subTest(trial=trial): + self.assertEqual(build_rolling(rows, decisions, specs), reference_rolling(rows, decisions, specs)) + + def test_input_order_does_not_change_features_or_lineage(self): + rows, decisions, specs = decode_rolling(rolling_demo_inputs()) + self.assertEqual(build_rolling(rows, decisions, specs), build_rolling(rows[::-1], decisions[::-1], specs[::-1])) + + +class RollingBundleTests(unittest.TestCase): + def rehash(self, root, filename): + manifest = read_json(root/"manifest.json") + manifest["files"][filename] = sha256((root/filename).read_bytes()).hexdigest() + (root/"manifest.json").write_text(canonical(manifest)) + + def test_full_bundle_replays_and_is_deterministic(self): + self.assertEqual(rolling_artifacts(rolling_demo_inputs()), rolling_artifacts(rolling_demo_inputs())) + with tempfile.TemporaryDirectory() as directory: + write_rolling_bundle(rolling_demo_inputs(), directory) + self.assertTrue(verify_rolling_bundle(directory)["verified"]) + + def test_rehashed_feature_fabrication_is_rejected(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + write_rolling_bundle(rolling_demo_inputs(), root) + content = (root/"features.csv").read_text() + self.assertIn("1500.0", content) + (root/"features.csv").write_text(content.replace("1500.0", "1501.0", 1)) + self.rehash(root, "features.csv") + with self.assertRaisesRegex(ValueError, "semantic replay"): + verify_rolling_bundle(root) + + def test_rehashed_member_fabrication_is_rejected(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + write_rolling_bundle(rolling_demo_inputs(), root) + content = (root/"members.csv").read_text() + (root/"members.csv").write_text(content.replace("cash-v1", "cash-v2", 1)) + self.rehash(root, "members.csv") + with self.assertRaisesRegex(ValueError, "semantic replay"): + verify_rolling_bundle(root) + + def test_rehashed_summary_fabrication_is_rejected(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + write_rolling_bundle(rolling_demo_inputs(), root) + summary = read_json(root/"summary.json") + summary["member_rows"] += 1 + (root/"summary.json").write_text(canonical(summary)) + self.rehash(root, "summary.json") + with self.assertRaisesRegex(ValueError, "semantic replay"): + verify_rolling_bundle(root) + + def test_empty_membership_still_exports_and_replays(self): + data = rolling_demo_inputs() + data["observations"] = [] + with tempfile.TemporaryDirectory() as directory: + write_rolling_bundle(data, directory) + self.assertEqual(len((Path(directory)/"members.csv").read_text().splitlines()), 1) + self.assertTrue(verify_rolling_bundle(directory)["verified"]) + + def test_identifiers_are_spreadsheet_safe_and_html_escaped(self): + data = rolling_demo_inputs() + malicious = '=SUM(1)' + data["decisions"][0]["decision_id"] = malicious + artifacts = rolling_artifacts(data) + self.assertIn("'=SUM", artifacts["features.csv"]) + self.assertNotIn('', artifacts["report.html"]) + self.assertIn("</script>", artifacts["report.html"]) + + def test_cli_demo_custom_input_and_verifier(self): + with tempfile.TemporaryDirectory() as directory, contextlib.redirect_stdout(io.StringIO()): + root = Path(directory) + self.assertEqual(main(["rolling-demo", "--out", str(root/"a")]), 0) + self.assertEqual(main(["aggregate", "--inputs", str(root/"a/inputs.json"), "--out", str(root/"b")]), 0) + self.assertEqual(main(["verify-rolling", "--out", str(root/"b")]), 0) + + def test_strict_input_schema_rejects_unknown_fields(self): + data = rolling_demo_inputs() + data["specs"] = [] + with self.assertRaises(ValueError): + decode_rolling(data)