"""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 @dataclass(frozen=True) class TradeTick: occurred_at: datetime price: float size: float bid: float | None = None ask: float | None = None serial: int | None = None @dataclass(frozen=True) 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) @dataclass(frozen=True) 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()} @dataclass(frozen=True) 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 @property 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 ) @property 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} @dataclass(frozen=True) class ChipFlowSnapshot: symbol: str source: str as_of: str available: bool windows: tuple[ChipWindow, ...] = () missing_reason: str | None = None shadow_metadata: ChipFlowShadowMetadata | None = None @property def status(self) -> str: return self.shadow_metadata.status if self.shadow_metadata is not None else "unconfirmed" @property 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