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