Cai
2026-08-29 3608f74e165a7118a602ce4b3d97d805e7912cb1
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
from __future__ import annotations
 
from typing import Any
 
from .cache import canonical_bytes, sha256_bytes
 
 
class CoreInputError(ValueError):
    pass
 
 
FINANCIAL_LABELS = ("annual", "current_cumulative", "prior_year_same_period")
FINANCIAL_VALUES = ("revenue", "attributable_profit", "deduct_profit", "cfo", "capex")
BALANCE_VALUES = (
    "equity",
    "cash_available",
    "non_operating_financial_assets",
    "interest_bearing_debt",
    "minority_interest",
)
CORE_VALUE_PATHS = (
    "market.price",
    "market.diluted_shares",
    "market.platform_market_cap",
    *(f"financials.{label}.{field}" for label in FINANCIAL_LABELS for field in FINANCIAL_VALUES),
    *(f"balance_sheet.{field}" for field in BALANCE_VALUES),
)
 
 
def protected_projection(snapshot: dict[str, Any]) -> dict[str, Any]:
    return {
        "meta": snapshot.get("meta"),
        "market": snapshot.get("market"),
        "financials": snapshot.get("financials"),
        "balance_sheet": snapshot.get("balance_sheet"),
        "sources": snapshot.get("sources"),
    }
 
 
def protected_hash(snapshot: dict[str, Any]) -> str:
    return sha256_bytes(canonical_bytes(protected_projection(snapshot)))
 
 
def _a1_for_period(
    sources: list[dict[str, Any]], period: str, support: str
) -> dict[str, Any] | None:
    candidates = [
        source
        for source in sources
        if source.get("source_type") in {"cninfo", "szse", "sse", "bse", "archived_primary"}
        and source.get("period_end") == period
        and support in source.get("supports", [])
        and source.get("publish_date")
        and source.get("id")
    ]
    return max(candidates, key=lambda item: (item["publish_date"], item["id"])) if candidates else None
 
 
def _require_lineage(
    path: str,
    lineage: dict[str, Any],
    source_by_id: dict[str, dict[str, Any]],
    all_raw: set[str],
    as_of: str,
) -> None:
    missing = [
        key for key in ("source_id", "raw_hash", "publish_date", "data_date", "provider", "raw_field", "unit")
        if not lineage.get(key)
    ]
    if missing:
        raise CoreInputError(f"字段血缘缺项 {path}: {','.join(missing)}")
    if lineage["source_id"] not in source_by_id:
        raise CoreInputError(f"字段血缘 source_id 不存在 {path}: {lineage['source_id']}")
    structured = source_by_id[lineage["source_id"]]
    if structured.get("period_end") != lineage["data_date"]:
        raise CoreInputError(f"字段血缘 data_date/source period 不一致:{path}")
    if structured.get("publish_date") != lineage["publish_date"]:
        raise CoreInputError(f"字段血缘 publish_date/source 不一致:{path}")
    if not any(
        path == support or path.startswith(f"{support}.")
        for support in structured.get("supports", [])
    ):
        raise CoreInputError(f"字段血缘结构化来源不支持字段:{path}")
    if lineage["raw_hash"] not in all_raw or len(lineage["raw_hash"]) != 64:
        raise CoreInputError(f"字段血缘 raw hash 不存在 {path}")
    if lineage["publish_date"] > as_of or lineage["data_date"] > as_of:
        raise CoreInputError(f"E_ASOF_VIOLATION:{path}")
 
 
def build_data_snapshot(
    ticker: str, as_of: str, providers: dict[str, dict[str, Any]]
) -> tuple[dict[str, Any], dict[str, Any]]:
    gaps = [gap for result in providers.values() for gap in result.get("gaps", [])]
    blocking = [gap for gap in gaps if gap.get("blocking")]
    if blocking:
        raise CoreInputError(";".join(gap["reason"] for gap in blocking))
    for required_provider in ("announcements", "market", "finance", "forecast"):
        if required_provider not in providers:
            raise CoreInputError(f"缺少 provider:{required_provider}")
    announcements = providers["announcements"]
    market = providers["market"]
    finance = providers["finance"]
    forecast = providers["forecast"]
    for key, container in (
        ("identity", announcements),
        ("market", market),
        ("financials", finance),
        ("balance_sheet", finance),
    ):
        if key not in container:
            raise CoreInputError(f"缺少核心字段组:{key}")
 
    sources: list[dict[str, Any]] = []
    seen: set[str] = set()
    for result in providers.values():
        for source in result.get("sources", []):
            source_id = source.get("id")
            if not source_id:
                raise CoreInputError("来源缺少 source_id")
            if not source.get("publish_date"):
                raise CoreInputError(f"来源缺少 publish_date:{source_id}")
            if source["publish_date"] > as_of:
                raise CoreInputError(f"E_ASOF_VIOLATION:{source_id}")
            if source.get("period_end") and source["period_end"] > as_of:
                raise CoreInputError(f"E_ASOF_VIOLATION:{source_id}")
            if source_id not in seen:
                seen.add(source_id)
                sources.append(source)
    source_by_id = {source["id"]: source for source in sources}
    all_raw = {
        digest
        for result in providers.values()
        for digest in result.get("raw_artifact_hashes", [])
    }
 
    financials = finance["financials"]
    current_period = financials["current_cumulative"]["period_end"]
    annual_period = financials["annual"]["period_end"]
    annual_a1 = _a1_for_period(sources, annual_period, "financials.annual")
    current_a1 = _a1_for_period(sources, current_period, "financials.current_cumulative")
    comparative_a1 = _a1_for_period(
        sources, current_period, "financials.prior_year_same_period_comparative"
    )
    balance_a1 = _a1_for_period(sources, current_period, "balance_sheet")
    if annual_a1 is None or current_a1 is None or comparative_a1 is None or balance_a1 is None:
        missing = []
        if annual_a1 is None:
            missing.append(annual_period)
        if current_a1 is None:
            missing.append(current_period)
        if comparative_a1 is None:
            missing.append(f"{current_period}:comparative")
        if balance_a1 is None:
            missing.append(f"{current_period}:balance")
        raise CoreInputError(f"核心财务缺少 A1 报告期回链:{missing}")
    announcement_hashes = announcements.get("raw_artifact_hashes", [])
    if not announcement_hashes:
        raise CoreInputError("A1 公告索引缺少 raw hash")
    a1_raw_hash = announcement_hashes[-1]
 
    field_lineage: dict[str, dict[str, Any]] = {}
    for result in providers.values():
        for path, item in result.get("field_lineage", {}).items():
            if path in field_lineage:
                raise CoreInputError(f"字段血缘重复:{path}")
            field_lineage[path] = dict(item)
    for path in CORE_VALUE_PATHS:
        if path not in field_lineage:
            raise CoreInputError(f"核心字段缺少血缘:{path}")
        _require_lineage(path, field_lineage[path], source_by_id, all_raw, as_of)
        if path.startswith("financials.") or path.startswith("balance_sheet."):
            if path.startswith("financials.annual."):
                a1 = annual_a1
                a1_support = "financials.annual"
            elif path.startswith("financials.prior_year_same_period."):
                a1 = comparative_a1
                a1_support = "financials.prior_year_same_period_comparative"
                relation = field_lineage[path].get("comparison_relation")
                expected_relation = {
                    "type": "same_response_comparative_row",
                    "current_period_end": current_period,
                    "comparison_period_end": financials["prior_year_same_period"]["period_end"],
                    "shared_publish_date": a1["publish_date"],
                }
                if relation != expected_relation:
                    raise CoreInputError(f"上年同期比较关系无效:{path}")
                if field_lineage[path]["publish_date"] != a1["publish_date"]:
                    raise CoreInputError(f"上年同期不是同一 A1 报告比较列:{path}")
            elif path.startswith("balance_sheet."):
                a1 = balance_a1
                a1_support = "balance_sheet"
            else:
                a1 = current_a1
                a1_support = "financials.current_cumulative"
            if a1_support not in a1.get("supports", []):
                raise CoreInputError(f"A1 未声明字段支持:{path}")
            field_lineage[path]["a1_source_id"] = a1["id"]
            field_lineage[path]["a1_publish_date"] = a1["publish_date"]
            field_lineage[path]["a1_period_end"] = a1["period_end"]
            field_lineage[path]["a1_raw_hash"] = a1_raw_hash
            field_lineage[path]["a1_support"] = a1_support
            if a1_raw_hash not in all_raw or a1["publish_date"] > as_of:
                raise CoreInputError(f"核心字段 A1 血缘无效:{path}")
 
    identity = announcements["identity"]
    snapshot = {
        "snapshot_id": f"SNAPSHOT-{ticker.split('.')[0]}-{as_of.replace('-', '')}-V2",
        "meta": {
            "company": identity["company"],
            "code": ticker,
            "market": identity["market"],
            "as_of_date": as_of,
            "currency": "CNY",
            "report_period_end": current_period,
            "latest_operating_info_date": max(source["publish_date"] for source in sources),
        },
        "market": market["market"],
        "financials": financials,
        "balance_sheet": finance["balance_sheet"],
        "institutions": forecast.get("institutions", {"coverage_status": "gap", "forecasts": []}),
        "sources": sources,
    }
    report = {
        "schema_version": 2,
        "ticker": ticker,
        "as_of": as_of,
        "field_lineage": field_lineage,
        "raw_hashes": sorted(all_raw),
        "gaps": gaps,
        "as_of_safe": True,
        "mechanical_protected_sha256": protected_hash(snapshot),
    }
    return snapshot, report