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
|