Cai
2026-08-20 f046e00dff100e7510089e4e568f2165ca8a199e
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
from __future__ import annotations
 
import time
from concurrent.futures import Future, ThreadPoolExecutor, wait
from typing import Any, Callable
 
from .cache import BlobIntegrityError
from .http_client import HttpClient
from .providers import acquire_announcements, acquire_finance, acquire_forecast, acquire_market
 
 
def acquire_all(
    client: HttpClient,
    ticker: str,
    as_of: str,
    baseline_results: dict[str, dict[str, Any]] | None = None,
) -> dict[str, dict[str, Any]]:
    jobs: dict[str, Callable[[HttpClient, str, str], dict[str, Any]]] = {
        "announcements": acquire_announcements,
        "market": acquire_market,
        "finance": acquire_finance,
        "forecast": acquire_forecast,
    }
    baseline_results = baseline_results or {}
    executor = ThreadPoolExecutor(max_workers=4, thread_name_prefix="valuation-v2")
    futures: dict[Future[dict[str, Any]], str] = {
        executor.submit(func, client, ticker, as_of): name
        for name, func in jobs.items()
        if name not in baseline_results
    }
    results: dict[str, dict[str, Any]] = dict(baseline_results)
    try:
        remaining = client.remaining()
        done, pending = wait(futures, timeout=remaining)
        for future in done:
            name = futures[future]
            try:
                results[name] = future.result()
            except BaseException as exc:
                blocking = name != "forecast"
                results[name] = {
                    "provider_id": name,
                    "adapter_version": "1.0.0",
                    "status": "BLOCKED" if blocking else "GAP",
                    "fetched_at": time.time(),
                    "as_of_date": as_of,
                    "records": [],
                    "sources": [],
                    "gaps": [
                        {
                            "gap_id": f"E_{name.upper()}_FAILED" if blocking else "W_FORECAST_FAILED",
                            "provider": name,
                            "field": name,
                            "reason": f"{type(exc).__name__}: {exc}",
                            "impact": "核心数据不可用" if blocking else "机构预测不可用",
                            "blocking": blocking,
                            "budget_used_seconds": None,
                            "manual_action": "核对登记来源或稍后重试",
                        }
                    ],
                    "warnings": [],
                    "raw_artifact_hashes": [],
                    "request_telemetry": [],
                    "cache_integrity_failure": isinstance(exc, BlobIntegrityError),
                }
        if pending and hasattr(client, "cancel"):
            client.cancel()
            finished_after_cancel, still_pending = wait(pending, timeout=1.0)
            pending = still_pending
            for future in finished_after_cancel:
                name = futures[future]
                try:
                    results[name] = future.result()
                except BaseException:
                    pending.add(future)
        for future in pending:
            future.cancel()
            name = futures[future]
            blocking = name != "forecast"
            results[name] = {
                "provider_id": name,
                "adapter_version": "1.0.0",
                "status": "BLOCKED" if blocking else "GAP",
                "fetched_at": time.time(),
                "as_of_date": as_of,
                "records": [],
                "sources": [],
                "gaps": [
                    {
                        "gap_id": "E_NETWORK_DEADLINE" if blocking else "W_NETWORK_DEADLINE",
                        "provider": name,
                        "field": name,
                        "reason": "90 秒网络截止已到",
                        "impact": "核心数据不可用" if blocking else "机构预测不可用",
                        "blocking": blocking,
                        "budget_used_seconds": 90.0,
                        "manual_action": "稍后重试",
                    }
                ],
                "warnings": [],
                "raw_artifact_hashes": [],
                "request_telemetry": [],
                "cache_integrity_failure": False,
            }
    finally:
        executor.shutdown(wait=False, cancel_futures=True)
    kinds = {
        "announcements": {"stock_identity", "announcement_index"},
        "market": {"market_close", "shares_market_cap"},
        "finance": {"finance_main", "finance_income", "finance_balance", "finance_cashflow"},
        "forecast": {"forecast_summary", "forecast_detail"},
    }
    telemetry = getattr(client, "telemetry", [])
    for name, result in results.items():
        data_kinds = set(result.get("data_kinds", {}))
        baseline_kinds = set(result.get("baseline_data_kinds", []))
        result["baseline_reused"] = bool(data_kinds) and data_kinds <= baseline_kinds
        if result.get("status") not in {"BLOCKED", "GAP"}:
            continue
        partial = [item for item in telemetry if item.get("data_kind") in kinds[name]]
        result["request_telemetry"] = partial
        result["raw_artifact_hashes"] = sorted(
            {item["raw_hash"] for item in partial if item.get("raw_hash")}
        )
    return results