DockerSpace / data /realtime_chip_adapter.py
DennisChan0909's picture
Backup current stock predictor strategies
ee37d63
Raw
History Blame Contribute Delete
15.7 kB
"""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