Spaces:
Sleeping
Sleeping
| """Provider-neutral adapter for explicit intraday big-player trade flow. | |
| The predictor must not infer broker-only chip data from OHLCV. This module | |
| keeps the external dependency behind a small interface so an MCP server, | |
| licensed feed, or mock fixture can provide the same normalized snapshot. | |
| """ | |
| from __future__ import annotations | |
| import json | |
| import os | |
| from dataclasses import asdict, dataclass | |
| from datetime import datetime, time, timezone | |
| from typing import Any, Callable, Literal, Protocol | |
| from urllib.request import Request, urlopen | |
| from zoneinfo import ZoneInfo | |
| TAIPEI = ZoneInfo("Asia/Taipei") | |
| LARGE_ORDER_MIN_TWD = 1_000_000.0 | |
| LARGE_ORDER_MAX_TWD = 4_000_000.0 | |
| EXTRA_LARGE_ORDER_MIN_TWD = 4_000_000.0 | |
| SHADOW_SCHEMA_VERSION = "realtime-chip-shadow.v1" | |
| PROVIDER_NEUTRAL_CONTRACT_VERSION = "explicit-trade-ticks.v1" | |
| FUGLE_PROVIDER_CONTRACT_VERSION = "fugle-marketdata-v1.0-intraday-trades.v1" | |
| FUGLE_MAX_TRADE_COUNT = 5_000 | |
| STALE_AFTER_SECONDS = 300.0 | |
| SEVERE_UNCLASSIFIED_RATIO = 0.25 | |
| class TradeTick: | |
| occurred_at: datetime | |
| price: float | |
| size: float | |
| bid: float | None = None | |
| ask: float | None = None | |
| serial: int | None = None | |
| class ChipFlowQuality: | |
| parse_drop_count: int = 0 | |
| duplicate_count: int = 0 | |
| serial_gap_count: int = 0 | |
| truncated: bool = False | |
| stale_seconds: float | None = None | |
| unclassified_ratio: float = 0.0 | |
| def to_dict(self) -> dict[str, Any]: | |
| return asdict(self) | |
| class ChipFlowShadowMetadata: | |
| schema_version: str | |
| provider_contract_version: str | |
| market_date: str | |
| captured_at: str | |
| status: Literal["confirmed", "partial", "unconfirmed", "invalid"] | |
| eligible_for_strict_gate: bool | |
| quality: ChipFlowQuality | |
| def to_dict(self) -> dict[str, Any]: | |
| return {**asdict(self), "quality": self.quality.to_dict()} | |
| class ChipWindow: | |
| window: str | |
| buy_big_order_1m_4m_twd: float | |
| sell_big_order_1m_4m_twd: float | |
| buy_super_order_gt_4m_twd: float | |
| sell_super_order_gt_4m_twd: float | |
| unclassified_twd: float | |
| tick_count: int = 0 | |
| complete: bool = False | |
| def net_twd(self) -> float: | |
| return ( | |
| self.buy_big_order_1m_4m_twd | |
| + self.buy_super_order_gt_4m_twd | |
| - self.sell_big_order_1m_4m_twd | |
| - self.sell_super_order_gt_4m_twd | |
| ) | |
| def sentiment(self) -> str: | |
| return "red" if self.net_twd > 0 else "green" if self.net_twd < 0 else "neutral" | |
| def to_dict(self) -> dict[str, Any]: | |
| return {**asdict(self), "net_twd": self.net_twd, "sentiment": self.sentiment} | |
| class ChipFlowSnapshot: | |
| symbol: str | |
| source: str | |
| as_of: str | |
| available: bool | |
| windows: tuple[ChipWindow, ...] = () | |
| missing_reason: str | None = None | |
| shadow_metadata: ChipFlowShadowMetadata | None = None | |
| def status(self) -> str: | |
| return self.shadow_metadata.status if self.shadow_metadata is not None else "unconfirmed" | |
| def eligible_for_strict_gate(self) -> bool: | |
| return self.shadow_metadata.eligible_for_strict_gate if self.shadow_metadata is not None else False | |
| def to_dict(self) -> dict[str, Any]: | |
| return { | |
| "symbol": self.symbol, | |
| "source": self.source, | |
| "as_of": self.as_of, | |
| "available": self.available, | |
| "last_two_60k": [window.to_dict() for window in self.windows], | |
| "missing_reason": self.missing_reason, | |
| "shadow_metadata": (self.shadow_metadata or _legacy_shadow_metadata(self)).to_dict(), | |
| } | |
| class RealtimeChipFlowAdapter(Protocol): | |
| """Transport boundary suitable for direct calls or an MCP wrapper.""" | |
| def snapshot(self, symbol: str, *, as_of: datetime | None = None) -> ChipFlowSnapshot: | |
| ... | |
| class UnavailableChipFlowAdapter: | |
| """Conservative fallback when no licensed tick provider is configured.""" | |
| def __init__(self, reason: str = "realtime chip provider is not configured") -> None: | |
| self.reason = reason | |
| def snapshot(self, symbol: str, *, as_of: datetime | None = None) -> ChipFlowSnapshot: | |
| now = _taipei_time(as_of) | |
| return ChipFlowSnapshot( | |
| symbol=symbol, | |
| source="unavailable", | |
| as_of=now.isoformat(), | |
| available=False, | |
| missing_reason=self.reason, | |
| shadow_metadata=_build_shadow_metadata(now, status="unconfirmed"), | |
| ) | |
| class MockChipFlowAdapter: | |
| """Deterministic adapter for tests and development without a provider key.""" | |
| def __init__(self, ticks_by_symbol: dict[str, list[TradeTick]], *, size_unit_shares: float = 1.0) -> None: | |
| self.ticks_by_symbol = ticks_by_symbol | |
| self.size_unit_shares = _positive_size_unit(size_unit_shares) | |
| def snapshot(self, symbol: str, *, as_of: datetime | None = None) -> ChipFlowSnapshot: | |
| now = _taipei_time(as_of) | |
| return aggregate_trade_ticks( | |
| symbol, | |
| self.ticks_by_symbol.get(symbol, []), | |
| source="mock", | |
| as_of=now, | |
| size_unit_shares=self.size_unit_shares, | |
| ) | |
| class FugleRestChipFlowAdapter: | |
| """Fugle REST adapter using the official intraday trades endpoint. | |
| Fugle documentation names the trade field ``size`` but the provider plan | |
| and unit contract must be confirmed before live use. Require an explicit | |
| ``size_unit_shares`` setting instead of guessing. | |
| """ | |
| BASE_URL = "https://api.fugle.tw/marketdata/v1.0/stock/intraday/trades" | |
| def __init__( | |
| self, | |
| api_key: str, | |
| *, | |
| size_unit_shares: float, | |
| requester: Callable[[Request], Any] | None = None, | |
| base_url: str = BASE_URL, | |
| ) -> None: | |
| if not api_key.strip(): | |
| raise ValueError("Fugle api_key is required") | |
| self.api_key = api_key.strip() | |
| self.size_unit_shares = _positive_size_unit(size_unit_shares) | |
| self.requester = requester or (lambda request: urlopen(request, timeout=15)) | |
| self.base_url = base_url.rstrip("/") | |
| def snapshot(self, symbol: str, *, as_of: datetime | None = None) -> ChipFlowSnapshot: | |
| now = _taipei_time(as_of) | |
| request = Request( | |
| f"{self.base_url}/{symbol}?limit=5000&sort=asc", | |
| headers={"X-API-KEY": self.api_key, "Accept": "application/json"}, | |
| ) | |
| with self.requester(request) as response: | |
| payload = json.loads(response.read().decode("utf-8")) | |
| raw_ticks = payload.get("data", []) | |
| malformed_payload = not isinstance(raw_ticks, list) | |
| if malformed_payload: | |
| raw_ticks = [] | |
| parsed_ticks = [_parse_fugle_tick(item) for item in raw_ticks] | |
| return aggregate_trade_ticks( | |
| symbol, | |
| [tick for tick in parsed_ticks if tick is not None], | |
| source="fugle-rest", | |
| as_of=now, | |
| size_unit_shares=self.size_unit_shares, | |
| provider_contract_version=FUGLE_PROVIDER_CONTRACT_VERSION, | |
| parse_drop_count=int(malformed_payload) + sum(tick is None for tick in parsed_ticks), | |
| truncated=_payload_is_truncated(payload, len(raw_ticks)), | |
| partial=bool(payload.get("partial") or payload.get("isPartial")), | |
| ) | |
| def build_realtime_chip_adapter(env: dict[str, str] | None = None) -> RealtimeChipFlowAdapter: | |
| """Build the configured adapter without silently guessing provider units.""" | |
| values = os.environ if env is None else env | |
| provider = values.get("REALTIME_CHIP_PROVIDER", "").strip().lower() | |
| if provider != "fugle": | |
| return UnavailableChipFlowAdapter("REALTIME_CHIP_PROVIDER is not configured") | |
| api_key = values.get("FUGLE_API_KEY", "").strip() | |
| unit = values.get("FUGLE_TRADE_SIZE_UNIT_SHARES", "").strip() | |
| if not api_key: | |
| return UnavailableChipFlowAdapter("FUGLE_API_KEY is missing") | |
| if not unit: | |
| return UnavailableChipFlowAdapter("FUGLE_TRADE_SIZE_UNIT_SHARES must be confirmed explicitly") | |
| try: | |
| return FugleRestChipFlowAdapter(api_key, size_unit_shares=float(unit)) | |
| except ValueError as exc: | |
| return UnavailableChipFlowAdapter(str(exc)) | |
| def aggregate_trade_ticks( | |
| symbol: str, | |
| ticks: list[TradeTick], | |
| *, | |
| source: str, | |
| as_of: datetime, | |
| size_unit_shares: float, | |
| provider_contract_version: str = PROVIDER_NEUTRAL_CONTRACT_VERSION, | |
| parse_drop_count: int = 0, | |
| truncated: bool = False, | |
| partial: bool = False, | |
| ) -> ChipFlowSnapshot: | |
| """Aggregate explicit trades into the final two Taiwan market windows.""" | |
| now = _taipei_time(as_of) | |
| unit = _positive_size_unit(size_unit_shares) | |
| deduplicated_ticks, duplicate_count = _deduplicate_ticks(ticks) | |
| relevant_ticks = [tick for tick in deduplicated_ticks if _taipei_time(tick.occurred_at).date() == now.date()] | |
| windows = [ | |
| ("11:30-12:30", time(11, 30), time(12, 30)), | |
| ("12:30-13:30", time(12, 30), time(13, 30)), | |
| ] | |
| buckets = [] | |
| for label, start, end in windows: | |
| values = { | |
| "buy_big_order_1m_4m_twd": 0.0, | |
| "sell_big_order_1m_4m_twd": 0.0, | |
| "buy_super_order_gt_4m_twd": 0.0, | |
| "sell_super_order_gt_4m_twd": 0.0, | |
| "unclassified_twd": 0.0, | |
| } | |
| tick_count = 0 | |
| for tick in relevant_ticks: | |
| occurred = _taipei_time(tick.occurred_at) | |
| if not start <= occurred.time() < end: | |
| continue | |
| tick_count += 1 | |
| amount = float(tick.price) * float(tick.size) * unit | |
| side = _trade_side(tick) | |
| if amount < LARGE_ORDER_MIN_TWD: | |
| continue | |
| if side is None: | |
| values["unclassified_twd"] += amount | |
| elif amount <= LARGE_ORDER_MAX_TWD: | |
| values[f"{side}_big_order_1m_4m_twd"] += amount | |
| else: | |
| values[f"{side}_super_order_gt_4m_twd"] += amount | |
| buckets.append(ChipWindow(window=label, tick_count=tick_count, complete=now.time() >= end, **values)) | |
| quality = ChipFlowQuality( | |
| parse_drop_count=max(0, int(parse_drop_count)), | |
| duplicate_count=duplicate_count, | |
| serial_gap_count=_serial_gap_count(deduplicated_ticks), | |
| truncated=bool(truncated), | |
| stale_seconds=_stale_seconds(relevant_ticks, now), | |
| unclassified_ratio=_unclassified_ratio(buckets), | |
| ) | |
| status = _shadow_status(now, relevant_ticks, quality, partial=partial) | |
| shadow_metadata = _build_shadow_metadata( | |
| now, | |
| provider_contract_version=provider_contract_version, | |
| status=status, | |
| quality=quality, | |
| ) | |
| return ChipFlowSnapshot( | |
| symbol=symbol, | |
| source=source, | |
| as_of=now.isoformat(), | |
| available=shadow_metadata.eligible_for_strict_gate, | |
| windows=tuple(buckets), | |
| shadow_metadata=shadow_metadata, | |
| ) | |
| def _trade_side(tick: TradeTick) -> str | None: | |
| if tick.ask is not None and tick.price >= tick.ask: | |
| return "buy" | |
| if tick.bid is not None and tick.price <= tick.bid: | |
| return "sell" | |
| return None | |
| def _parse_fugle_tick(item: dict[str, Any]) -> TradeTick | None: | |
| try: | |
| occurred = datetime.fromtimestamp(float(item["time"]) / 1_000_000.0, tz=timezone.utc) | |
| return TradeTick( | |
| occurred_at=occurred, | |
| price=float(item["price"]), | |
| size=float(item["size"]), | |
| bid=float(item["bid"]) if item.get("bid") is not None else None, | |
| ask=float(item["ask"]) if item.get("ask") is not None else None, | |
| serial=int(item["serial"]), | |
| ) | |
| except (AttributeError, KeyError, TypeError, ValueError): | |
| return None | |
| def _deduplicate_ticks(ticks: list[TradeTick]) -> tuple[list[TradeTick], int]: | |
| deduplicated = [] | |
| seen: set[tuple[Any, ...]] = set() | |
| duplicate_count = 0 | |
| for tick in ticks: | |
| key = ( | |
| ("serial", tick.serial) | |
| if tick.serial is not None | |
| else ("tick", _taipei_time(tick.occurred_at).isoformat(), tick.price, tick.size, tick.bid, tick.ask) | |
| ) | |
| if key in seen: | |
| duplicate_count += 1 | |
| continue | |
| seen.add(key) | |
| deduplicated.append(tick) | |
| return deduplicated, duplicate_count | |
| def _serial_gap_count(ticks: list[TradeTick]) -> int: | |
| serials = sorted({tick.serial for tick in ticks if tick.serial is not None}) | |
| return sum(max(0, current - previous - 1) for previous, current in zip(serials, serials[1:])) | |
| def _stale_seconds(ticks: list[TradeTick], captured_at: datetime) -> float | None: | |
| if not ticks: | |
| return None | |
| latest_tick = max(_taipei_time(tick.occurred_at) for tick in ticks) | |
| return max(0.0, (captured_at - latest_tick).total_seconds()) | |
| def _unclassified_ratio(windows: list[ChipWindow]) -> float: | |
| unclassified_twd = sum(window.unclassified_twd for window in windows) | |
| total_twd = unclassified_twd + sum( | |
| window.buy_big_order_1m_4m_twd | |
| + window.sell_big_order_1m_4m_twd | |
| + window.buy_super_order_gt_4m_twd | |
| + window.sell_super_order_gt_4m_twd | |
| for window in windows | |
| ) | |
| return unclassified_twd / total_twd if total_twd else 0.0 | |
| def _shadow_status( | |
| captured_at: datetime, | |
| ticks: list[TradeTick], | |
| quality: ChipFlowQuality, | |
| *, | |
| partial: bool, | |
| ) -> Literal["confirmed", "partial", "unconfirmed", "invalid"]: | |
| if ( | |
| quality.parse_drop_count | |
| or quality.duplicate_count | |
| or quality.serial_gap_count | |
| or quality.unclassified_ratio > SEVERE_UNCLASSIFIED_RATIO | |
| ): | |
| return "invalid" | |
| if not ticks: | |
| return "unconfirmed" | |
| if ( | |
| partial | |
| or quality.truncated | |
| or captured_at.time() < time(13, 30) | |
| or quality.stale_seconds is None | |
| or quality.stale_seconds > STALE_AFTER_SECONDS | |
| ): | |
| return "partial" | |
| return "confirmed" | |
| def _build_shadow_metadata( | |
| captured_at: datetime, | |
| *, | |
| provider_contract_version: str = PROVIDER_NEUTRAL_CONTRACT_VERSION, | |
| status: Literal["confirmed", "partial", "unconfirmed", "invalid"], | |
| quality: ChipFlowQuality | None = None, | |
| ) -> ChipFlowShadowMetadata: | |
| return ChipFlowShadowMetadata( | |
| schema_version=SHADOW_SCHEMA_VERSION, | |
| provider_contract_version=provider_contract_version, | |
| market_date=captured_at.date().isoformat(), | |
| captured_at=captured_at.isoformat(), | |
| status=status, | |
| eligible_for_strict_gate=status == "confirmed", | |
| quality=quality or ChipFlowQuality(), | |
| ) | |
| def _legacy_shadow_metadata(snapshot: ChipFlowSnapshot) -> ChipFlowShadowMetadata: | |
| try: | |
| captured_at = _taipei_time(datetime.fromisoformat(snapshot.as_of)) | |
| except ValueError: | |
| captured_at = _taipei_time(None) | |
| return _build_shadow_metadata(captured_at, status="unconfirmed") | |
| def _payload_is_truncated(payload: dict[str, Any], trade_count: int) -> bool: | |
| return bool( | |
| trade_count >= FUGLE_MAX_TRADE_COUNT | |
| or payload.get("truncated") | |
| or payload.get("isTruncated") | |
| or payload.get("hasNextPage") | |
| or payload.get("nextPage") | |
| or payload.get("nextPageToken") | |
| ) | |
| def _taipei_time(value: datetime | None) -> datetime: | |
| current = value or datetime.now(tz=TAIPEI) | |
| if current.tzinfo is None: | |
| current = current.replace(tzinfo=TAIPEI) | |
| return current.astimezone(TAIPEI) | |
| def _positive_size_unit(value: float) -> float: | |
| unit = float(value) | |
| if unit <= 0: | |
| raise ValueError("trade size unit must be positive") | |
| return unit | |