From 35b846f1d6c1661f99dd6a5fabe85953a48a349b Mon Sep 17 00:00:00 2001 From: "2569718930@qq.com" <2569718930@qq.com> Date: Tue, 10 Mar 2026 12:19:23 +0800 Subject: [PATCH] feat: Add Polymarket and Polygon wallet activity watchers and integrate them into the bot. --- .env.example | 16 + bot_listener.py | 3 + src/onchain/polygon_wallet_watcher.py | 260 +++++++++-- .../polymarket_wallet_activity_watcher.py | 413 ++++++++++++++++++ 4 files changed, 649 insertions(+), 43 deletions(-) create mode 100644 src/onchain/polymarket_wallet_activity_watcher.py diff --git a/.env.example b/.env.example index de4f2981..472f61bd 100644 --- a/.env.example +++ b/.env.example @@ -48,3 +48,19 @@ POLYGON_WALLET_WATCH_SEEN_TTL_SEC=604800 POLYGON_WALLET_WATCH_RPC_TIMEOUT_SEC=10 POLYGON_WALLET_WATCH_TX_BASE=https://polygonscan.com/tx POLYGON_WALLET_WATCH_ADDR_BASE=https://polygonscan.com/address +POLYGON_WALLET_WATCH_POLYMARKET_ONLY=true +POLYGON_WALLET_WATCH_INCLUDE_DEFAULT_PM_CONTRACTS=true +# Optional custom Polymarket contracts, format: LABEL:0x...,LABEL2:0x... +POLYGON_WALLET_WATCH_POLYMARKET_CONTRACTS= + +# Polymarket Wallet Activity Watcher (all markets, not weather-only) +POLYMARKET_WALLET_ACTIVITY_ENABLED=false +POLYMARKET_WALLET_ACTIVITY_USERS=0x0000000000000000000000000000000000000000 +POLYMARKET_WALLET_ACTIVITY_DATA_API_URL=https://data-api.polymarket.com +POLYMARKET_WALLET_ACTIVITY_INTERVAL_SEC=20 +POLYMARKET_WALLET_ACTIVITY_TIMEOUT_SEC=10 +POLYMARKET_WALLET_ACTIVITY_MIN_SIZE_ABS=0.001 +POLYMARKET_WALLET_ACTIVITY_MIN_SIZE_DELTA=0.001 +POLYMARKET_WALLET_ACTIVITY_MAX_CHANGES_PER_MSG=5 +POLYMARKET_WALLET_ACTIVITY_NOTIFY_CLOSED=true +POLYMARKET_WALLET_ACTIVITY_BOOTSTRAP_ALERT=false diff --git a/bot_listener.py b/bot_listener.py index 0117a08a..b2d65c7d 100644 --- a/bot_listener.py +++ b/bot_listener.py @@ -12,6 +12,7 @@ if project_root not in sys.path: from src.utils.config_loader import load_config # type: ignore # noqa: E402 from src.utils.telegram_push import start_trade_alert_push_loop # type: ignore # noqa: E402 from src.onchain.polygon_wallet_watcher import start_polygon_wallet_watch_loop # type: ignore # noqa: E402 +from src.onchain.polymarket_wallet_activity_watcher import start_polymarket_wallet_activity_loop # type: ignore # noqa: E402 from src.data_collection.weather_sources import WeatherDataCollector # type: ignore # noqa: E402 from src.data_collection.city_risk_profiles import get_city_risk_profile # type: ignore # noqa: E402 from src.analysis.deb_algorithm import calculate_dynamic_weights, update_daily_record # noqa: E402 @@ -44,6 +45,7 @@ def start_bot(): weather = WeatherDataCollector(config) start_trade_alert_push_loop(bot, config) start_polygon_wallet_watch_loop(bot) + start_polymarket_wallet_activity_loop(bot) def _display_name(user) -> str: return user.username or user.first_name or f"User_{user.id}" @@ -813,3 +815,4 @@ def start_bot(): if __name__ == "__main__": start_bot() + diff --git a/src/onchain/polygon_wallet_watcher.py b/src/onchain/polygon_wallet_watcher.py index b43b59a8..2fcbf82d 100644 --- a/src/onchain/polygon_wallet_watcher.py +++ b/src/onchain/polygon_wallet_watcher.py @@ -10,6 +10,7 @@ from loguru import logger from web3 import Web3 TRANSFER_TOPIC = "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef" +APPROVAL_TOPIC = "0x8c5be1e5ebec7d5bd14f71427d1e84f3dd0314c0f7b2291e5b200ac8c7c3b925" ERC20_ABI = [ { "constant": True, @@ -27,6 +28,16 @@ ERC20_ABI = [ }, ] +# Source: Polymarket official developer docs (Polygon contract addresses) +# https://docs.polymarket.com/developers/market-makers/setup +DEFAULT_POLYMARKET_CONTRACTS: Dict[str, str] = { + "USDC.e": "0x2791Bca1f2de4661ED88A30C99A7a9449Aa84174", + "CTF": "0x4d97dcd97ec945f40cf65f87097ace5ea0476045", + "CTF_EXCHANGE": "0x4bFb41d5B3570DeFd03C39a9A4D8dE6Bd8B8982E", + "NEG_RISK_CTF_EXCHANGE": "0xC5d563A36AE78145C45a50134d48A1215220f80a", + "NEG_RISK_ADAPTER": "0xd91E80cF2E7be2e162c6513ceD06f1dD0dA35296", +} + def _env_bool(name: str, default: bool) -> bool: raw = os.getenv(name) @@ -91,17 +102,72 @@ def _cleanup_seen_tx(state: Dict[str, Any], now_ts: int, keep_sec: int) -> None: seen.pop(tx_hash, None) +def _normalize_addr(value: Any) -> str: + if value is None: + return "" + text = str(value).strip().lower() + if text.startswith("0x") and len(text) == 42: + return text + return "" + + def _parse_addresses(raw: Optional[str]) -> Set[str]: out: Set[str] = set() if not raw: return out for part in raw.split(","): - addr = part.strip().lower() - if addr and addr.startswith("0x") and len(addr) == 42: + addr = _normalize_addr(part) + if addr: out.add(addr) return out +def _parse_polymarket_contracts(raw: Optional[str]) -> Dict[str, str]: + """ + Parse env like: + - "0xabc...,0xdef..." + - "CTF:0xabc...,EXCHANGE:0xdef..." + """ + result: Dict[str, str] = {} + if not raw: + return result + for part in raw.split(","): + segment = str(part).strip() + if not segment: + continue + label = "CUSTOM_PM" + address_part = segment + if ":" in segment: + maybe_label, maybe_addr = segment.split(":", 1) + maybe_addr_n = _normalize_addr(maybe_addr) + if maybe_addr_n: + label = (maybe_label or "CUSTOM_PM").strip() or "CUSTOM_PM" + address_part = maybe_addr_n + addr = _normalize_addr(address_part) + if not addr: + continue + if addr not in result: + result[addr] = label + return result + + +def _build_polymarket_contract_map() -> Dict[str, str]: + include_defaults = _env_bool("POLYGON_WALLET_WATCH_INCLUDE_DEFAULT_PM_CONTRACTS", True) + merged: Dict[str, str] = {} + + if include_defaults: + for label, addr in DEFAULT_POLYMARKET_CONTRACTS.items(): + normalized = _normalize_addr(addr) + if normalized: + merged[normalized] = label + + custom = _parse_polymarket_contracts(os.getenv("POLYGON_WALLET_WATCH_POLYMARKET_CONTRACTS")) + for addr, label in custom.items(): + merged[addr] = label + + return merged + + def _polygon_scan_tx_url(tx_hash: str) -> str: base = os.getenv("POLYGON_WALLET_WATCH_TX_BASE") or "https://polygonscan.com/tx/" return f"{base.rstrip('/')}/{tx_hash}" @@ -123,69 +189,139 @@ def _format_matic(wei_value: int) -> str: return f"{matic.normalize():f}".rstrip("0").rstrip(".") +def _format_amount(amount: Decimal) -> str: + if amount == amount.to_integral_value(): + return str(int(amount)) + return f"{amount.normalize():f}".rstrip("0").rstrip(".") + + def _safe_lower(value: Any) -> str: if value is None: return "" return str(value).lower() -def _extract_transfer_events( +def _topic_to_addr(topic: Any) -> str: + try: + return _normalize_addr("0x" + topic.hex()[-40:]) + except Exception: + return "" + + +def _get_token_meta( + w3: Web3, + token_addr: str, + token_meta_cache: Dict[str, Tuple[str, int]], +) -> Tuple[str, int]: + symbol, decimals = token_meta_cache.get(token_addr, ("ERC20", 18)) + if token_addr in token_meta_cache: + return symbol, decimals + + try: + token = w3.eth.contract(address=Web3.to_checksum_address(token_addr), abi=ERC20_ABI) + symbol_raw = token.functions.symbol().call() + decimals_raw = token.functions.decimals().call() + symbol = str(symbol_raw or "ERC20") + decimals = int(decimals_raw) + except Exception: + symbol = "ERC20" + decimals = 18 + + token_meta_cache[token_addr] = (symbol, decimals) + return symbol, decimals + + +def _extract_receipt_signals( w3: Web3, receipt: Any, watch_set: Set[str], + pm_contracts: Dict[str, str], token_meta_cache: Dict[str, Tuple[str, int]], -) -> List[str]: - lines: List[str] = [] +) -> Dict[str, Any]: + transfer_lines: List[str] = [] + approval_lines: List[str] = [] + touched_labels: Set[str] = set() + pm_hit = False + for log in receipt.logs or []: try: - if not log.topics or len(log.topics) < 3: - continue - topic0 = log.topics[0].hex().lower() - if topic0 != TRANSFER_TOPIC: + log_addr = _normalize_addr(log.address) + if log_addr in pm_contracts: + pm_hit = True + touched_labels.add(pm_contracts[log_addr]) + + topics = log.topics or [] + if not topics: continue + topic0 = topics[0].hex().lower() - from_addr = "0x" + log.topics[1].hex()[-40:] - to_addr = "0x" + log.topics[2].hex()[-40:] - from_addr = from_addr.lower() - to_addr = to_addr.lower() - if from_addr not in watch_set and to_addr not in watch_set: - continue + if topic0 == TRANSFER_TOPIC and len(topics) >= 3: + from_addr = _topic_to_addr(topics[1]) + to_addr = _topic_to_addr(topics[2]) + if from_addr not in watch_set and to_addr not in watch_set: + continue - token_addr = _safe_lower(log.address) - symbol, decimals = token_meta_cache.get(token_addr, ("ERC20", 18)) - if token_addr not in token_meta_cache: - try: - token = w3.eth.contract(address=Web3.to_checksum_address(token_addr), abi=ERC20_ABI) - symbol_raw = token.functions.symbol().call() - decimals_raw = token.functions.decimals().call() - symbol = str(symbol_raw or "ERC20") - decimals = int(decimals_raw) - except Exception: - symbol = "ERC20" - decimals = 18 - token_meta_cache[token_addr] = (symbol, decimals) + other_addr = to_addr if from_addr in watch_set else from_addr + other_label = pm_contracts.get(other_addr) + if other_label: + pm_hit = True + touched_labels.add(other_label) - amount_int = int(log.data.hex(), 16) if log.data else 0 - amount = Decimal(amount_int) / (Decimal(10) ** Decimal(max(decimals, 0))) - amount_txt = f"{amount.normalize():f}".rstrip("0").rstrip(".") if amount else "0" + symbol, decimals = _get_token_meta(w3, log_addr, token_meta_cache) + amount_int = int(log.data.hex(), 16) if log.data else 0 + amount = Decimal(amount_int) / (Decimal(10) ** Decimal(max(decimals, 0))) - direction = "IN" if to_addr in watch_set and from_addr not in watch_set else "OUT" - if from_addr in watch_set and to_addr in watch_set: - direction = "SELF" + if from_addr in watch_set and to_addr in watch_set: + direction = "SELF" + elif to_addr in watch_set: + direction = "IN" + else: + direction = "OUT" - lines.append( - f"- {direction} {symbol}: {amount_txt} ({_short(from_addr)} -> {_short(to_addr)})" - ) + # Keep transfer line only when it is clearly Polymarket related. + if other_label or log_addr in pm_contracts: + target = other_label or pm_contracts.get(log_addr) or _short(other_addr) + transfer_lines.append( + f"- {direction} {symbol}: {_format_amount(amount)} (对手: {target})" + ) + + if topic0 == APPROVAL_TOPIC and len(topics) >= 3: + owner = _topic_to_addr(topics[1]) + spender = _topic_to_addr(topics[2]) + if owner not in watch_set: + continue + spender_label = pm_contracts.get(spender) + if not spender_label: + continue + + pm_hit = True + touched_labels.add(spender_label) + + symbol, decimals = _get_token_meta(w3, log_addr, token_meta_cache) + amount_int = int(log.data.hex(), 16) if log.data else 0 + amount = Decimal(amount_int) / (Decimal(10) ** Decimal(max(decimals, 0))) + approval_lines.append( + f"- APPROVE {symbol}: {_format_amount(amount)} -> {spender_label}" + ) except Exception: continue - return lines + + return { + "pm_hit": pm_hit, + "transfer_lines": transfer_lines, + "approval_lines": approval_lines, + "touched_labels": sorted(touched_labels), + } def _build_message( tx: Any, block_ts: int, matched_wallet: str, + touched_labels: List[str], transfer_lines: List[str], + approval_lines: List[str], + tx_to_label: Optional[str], ) -> str: tx_hash = tx["hash"].hex() from_addr = _safe_lower(tx.get("from")) @@ -202,21 +338,33 @@ def _build_message( matic_value = int(tx.get("value", 0) or 0) block_time = datetime.fromtimestamp(block_ts, tz=timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC") + selector = str(tx.get("input") or "")[:10] if tx.get("input") else "0x" lines = [ - "⛓ Polygon 钱包异动", + "⛓ Polymarket 钱包动作", f"钱包: {_short(matched_wallet)}", f"方向: {direction}", f"MATIC: {_format_matic(matic_value)}", f"区块: {tx.get('blockNumber')}", f"时间: {block_time}", - f"Tx: {tx_hash}", + f"方法选择器: {selector}", ] + if tx_to_label: + lines.append(f"直连合约: {tx_to_label}") + + if touched_labels: + lines.append(f"相关合约: {', '.join(touched_labels)}") + if transfer_lines: - lines.append("Token 转账:") + lines.append("Token 动作:") lines.extend(transfer_lines[:6]) + if approval_lines: + lines.append("授权动作:") + lines.extend(approval_lines[:4]) + + lines.append(f"Tx: {tx_hash}") lines.append(f"交易链接: {_polygon_scan_tx_url(tx_hash)}") lines.append(f"钱包链接: {_polygon_scan_addr_url(matched_wallet)}") return "\n".join(lines) @@ -227,6 +375,8 @@ def start_polygon_wallet_watch_loop(bot: Any) -> Optional[threading.Thread]: chat_id = os.getenv("TELEGRAM_CHAT_ID") rpc_url = os.getenv("POLYGON_RPC_URL") watch_set = _parse_addresses(os.getenv("POLYGON_WALLET_WATCH_ADDRESSES")) + polymarket_only = _env_bool("POLYGON_WALLET_WATCH_POLYMARKET_ONLY", True) + pm_contracts = _build_polymarket_contract_map() if not enabled: logger.info("polygon wallet watcher disabled") @@ -240,6 +390,9 @@ def start_polygon_wallet_watch_loop(bot: Any) -> Optional[threading.Thread]: if not watch_set: logger.warning("polygon wallet watcher skipped: POLYGON_WALLET_WATCH_ADDRESSES is empty") return None + if polymarket_only and not pm_contracts: + logger.warning("polygon wallet watcher skipped: no polymarket contracts configured") + return None poll_sec = max(3, _env_int("POLYGON_WALLET_WATCH_INTERVAL_SEC", 8)) confirmations = max(0, _env_int("POLYGON_WALLET_WATCH_CONFIRMATIONS", 2)) @@ -272,6 +425,7 @@ def start_polygon_wallet_watch_loop(bot: Any) -> Optional[threading.Thread]: logger.info( f"polygon wallet watcher started wallets={len(watch_set)} " + f"polymarket_only={polymarket_only} pm_contracts={len(pm_contracts)} " f"poll={poll_sec}s confirmations={confirmations} state_path={state_path}" ) @@ -312,29 +466,48 @@ def start_polygon_wallet_watch_loop(bot: Any) -> Optional[threading.Thread]: if tx_hash in state.get("seen_tx", {}): continue + tx_to = _normalize_addr(tx.get("to")) + tx_to_label = pm_contracts.get(tx_to) + pm_hit = bool(tx_to_label) transfer_lines: List[str] = [] + approval_lines: List[str] = [] + touched_labels: List[str] = [tx_to_label] if tx_to_label else [] + try: receipt = w3.eth.get_transaction_receipt(tx["hash"]) - transfer_lines = _extract_transfer_events( + parsed = _extract_receipt_signals( w3=w3, receipt=receipt, watch_set=watch_set, + pm_contracts=pm_contracts, token_meta_cache=token_meta_cache, ) + pm_hit = pm_hit or bool(parsed.get("pm_hit")) + transfer_lines = parsed.get("transfer_lines") or [] + approval_lines = parsed.get("approval_lines") or [] + touched = parsed.get("touched_labels") or [] + touched_labels = sorted(set(touched_labels + touched)) except Exception: transfer_lines = [] + approval_lines = [] + + if polymarket_only and not pm_hit: + continue message = _build_message( tx=tx, block_ts=block_ts, matched_wallet=matched_wallet, + touched_labels=touched_labels, transfer_lines=transfer_lines, + approval_lines=approval_lines, + tx_to_label=tx_to_label, ) bot.send_message(chat_id, message, disable_web_page_preview=True) state.setdefault("seen_tx", {})[tx_hash] = cycle_ts logger.info( f"polygon wallet alert pushed wallet={matched_wallet} " - f"tx={tx_hash} block={block_num}" + f"tx={tx_hash} block={block_num} polymarket={pm_hit}" ) state["last_scanned_block"] = block_num @@ -352,3 +525,4 @@ def start_polygon_wallet_watch_loop(bot: Any) -> Optional[threading.Thread]: ) thread.start() return thread + diff --git a/src/onchain/polymarket_wallet_activity_watcher.py b/src/onchain/polymarket_wallet_activity_watcher.py new file mode 100644 index 00000000..175cbdc9 --- /dev/null +++ b/src/onchain/polymarket_wallet_activity_watcher.py @@ -0,0 +1,413 @@ +import json +import os +import threading +import time +from datetime import datetime, timezone +from typing import Any, Dict, List, Optional, Tuple + +import requests +from loguru import logger + + +def _env_bool(name: str, default: bool) -> bool: + raw = os.getenv(name) + if raw is None: + return default + return raw.strip().lower() in {"1", "true", "yes", "on"} + + +def _env_int(name: str, default: int) -> int: + raw = os.getenv(name) + if raw is None: + return default + try: + return int(raw) + except Exception: + return default + + +def _env_float(name: str, default: float) -> float: + raw = os.getenv(name) + if raw is None: + return default + try: + return float(raw) + except Exception: + return default + + +def _safe_float(value: Any, default: float = 0.0) -> float: + try: + if value is None: + return default + return float(value) + except Exception: + return default + + +def _normalize_addr(value: Any) -> str: + if value is None: + return "" + text = str(value).strip().lower() + if text.startswith("0x") and len(text) == 42: + return text + return "" + + +def _short(addr: str, left: int = 6, right: int = 4) -> str: + if not addr: + return "unknown" + if len(addr) <= left + right + 2: + return addr + return f"{addr[:left + 2]}...{addr[-right:]}" + + +def _parse_addresses(raw: Optional[str]) -> List[str]: + out: List[str] = [] + if not raw: + return out + seen = set() + for part in raw.split(","): + addr = _normalize_addr(part) + if addr and addr not in seen: + out.append(addr) + seen.add(addr) + return out + + +def _state_file() -> str: + root = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + return os.path.join(root, "data", "polymarket_wallet_activity_state.json") + + +def _load_state(path: str) -> Dict[str, Any]: + if not os.path.exists(path): + return {"users": {}} + try: + with open(path, "r", encoding="utf-8") as fh: + data = json.load(fh) + if isinstance(data, dict): + data.setdefault("users", {}) + return data + except Exception as exc: + logger.warning(f"failed to load wallet activity state: {exc}") + return {"users": {}} + + +def _save_state(path: str, state: Dict[str, Any]) -> None: + os.makedirs(os.path.dirname(path), exist_ok=True) + tmp = f"{path}.tmp" + with open(tmp, "w", encoding="utf-8") as fh: + json.dump(state, fh, ensure_ascii=False, indent=2) + os.replace(tmp, path) + + +def _market_url(position: Dict[str, Any]) -> str: + slug = str(position.get("slug") or "").strip() + event_slug = str(position.get("event_slug") or "").strip() + + if slug and event_slug: + return f"https://polymarket.com/event/{event_slug}/{slug}" + if slug: + return f"https://polymarket.com/market/{slug}" + if event_slug: + return f"https://polymarket.com/event/{event_slug}" + return "" + + +def _position_key(position: Dict[str, Any]) -> str: + asset = str(position.get("asset") or "").strip().lower() + condition_id = str(position.get("condition_id") or "").strip().lower() + outcome = str(position.get("outcome") or "").strip().lower() + if asset: + return f"asset:{asset}" + return f"condition:{condition_id}|outcome:{outcome}" + + +def _normalize_position(row: Dict[str, Any]) -> Dict[str, Any]: + title = ( + str( + row.get("title") + or row.get("question") + or row.get("market") + or row.get("name") + or "" + ) + .strip() + ) + outcome = str(row.get("outcome") or row.get("side") or "").strip() + + return { + "proxy_wallet": _normalize_addr(row.get("proxyWallet") or row.get("proxy_wallet")), + "asset": str(row.get("asset") or row.get("assetId") or "").strip(), + "condition_id": str(row.get("conditionId") or row.get("condition_id") or "").strip(), + "title": title, + "slug": str(row.get("slug") or row.get("marketSlug") or "").strip(), + "event_slug": str(row.get("eventSlug") or row.get("event_slug") or "").strip(), + "outcome": outcome, + "size": _safe_float(row.get("size") or row.get("shares")), + "avg_price": _safe_float(row.get("avgPrice") or row.get("avg_price")), + "position_value": _safe_float(row.get("currentValue") or row.get("positionValue")), + "cash_pnl": _safe_float(row.get("cashPnl") or row.get("realizedPnl") or row.get("pnl")), + "percent_pnl": _safe_float( + row.get("percentPnl") + or row.get("pnlPercent") + or row.get("pnl_pct") + ), + "cur_price": _safe_float(row.get("curPrice") or row.get("lastPrice")), + "updated_at": int(time.time()), + } + + +def _fetch_positions( + session: requests.Session, + base_url: str, + user: str, + timeout_sec: int, +) -> List[Dict[str, Any]]: + url = f"{base_url.rstrip('/')}/positions" + resp = session.get(url, params={"user": user}, timeout=timeout_sec) + resp.raise_for_status() + data = resp.json() + + if isinstance(data, list): + return [row for row in data if isinstance(row, dict)] + + if isinstance(data, dict): + for key in ("positions", "data", "results", "items"): + maybe = data.get(key) + if isinstance(maybe, list): + return [row for row in maybe if isinstance(row, dict)] + + return [] + + +def _build_snapshot( + rows: List[Dict[str, Any]], + min_size_abs: float, +) -> Dict[str, Dict[str, Any]]: + snap: Dict[str, Dict[str, Any]] = {} + for row in rows: + pos = _normalize_position(row) + if abs(pos["size"]) < min_size_abs: + continue + key = _position_key(pos) + snap[key] = pos + return snap + + +def _diff_positions( + previous: Dict[str, Dict[str, Any]], + current: Dict[str, Dict[str, Any]], + min_size_delta: float, + notify_closed: bool, +) -> List[Tuple[str, Dict[str, Any]]]: + changes: List[Tuple[str, Dict[str, Any]]] = [] + + for key, now_pos in current.items(): + old_pos = previous.get(key) + if old_pos is None: + changes.append(("new", now_pos)) + continue + + size_delta = now_pos["size"] - old_pos.get("size", 0.0) + avg_delta = now_pos["avg_price"] - old_pos.get("avg_price", 0.0) + + if abs(size_delta) >= min_size_delta or abs(avg_delta) >= 1e-9: + merged = {**now_pos} + merged["size_delta"] = size_delta + merged["old_size"] = old_pos.get("size", 0.0) + merged["old_avg_price"] = old_pos.get("avg_price", 0.0) + changes.append(("update", merged)) + + if notify_closed: + for key, old_pos in previous.items(): + if key not in current: + changes.append(("closed", old_pos)) + + return changes + + +def _fmt_pct(value: float) -> str: + # Data API may return either ratio (0.12) or percent (12.0). + display = value * 100.0 if abs(value) <= 1.5 else value + return f"{display:.1f}%" + + +def _fmt_usd(value: float) -> str: + return f"${value:.2f}" + + +def _fmt_price(value: float) -> str: + return f"{value:.3f}" + + +def _format_change_block( + change_type: str, + wallet: str, + pos: Dict[str, Any], + now_utc: str, +) -> str: + title = pos.get("title") or "Unknown market" + outcome = pos.get("outcome") or "Unknown" + market_url = _market_url(pos) + + lines: List[str] = [] + if change_type == "new": + lines.append("🆕 New Position") + elif change_type == "update": + lines.append("🔄 Position Update") + else: + lines.append("❌ Position Closed") + + lines.append(f"Wallet: {_short(wallet)}") + if market_url: + lines.append(f"Market: {title} ({market_url})") + else: + lines.append(f"Market: {title}") + + lines.append(f"Outcome: {outcome}") + + if change_type == "update": + old_size = _safe_float(pos.get("old_size")) + now_size = _safe_float(pos.get("size")) + delta = _safe_float(pos.get("size_delta")) + lines.append(f"Size: {old_size:.3f} -> {now_size:.3f} (Δ {delta:+.3f})") + elif change_type == "closed": + lines.append(f"Size: 0.000 (was {_safe_float(pos.get('size')):.3f})") + else: + lines.append(f"Size: {_safe_float(pos.get('size')):.3f}") + + lines.append(f"Avg Price: {_fmt_price(_safe_float(pos.get('avg_price')))}") + lines.append(f"Position Value: {_fmt_usd(_safe_float(pos.get('position_value')))}") + + pnl = _safe_float(pos.get("cash_pnl")) + pnl_pct = _safe_float(pos.get("percent_pnl")) + lines.append(f"PnL: {_fmt_usd(pnl)} ({_fmt_pct(pnl_pct)})") + lines.append(f"Time: {now_utc}") + return "\n".join(lines) + + +def _build_message( + wallet: str, + changes: List[Tuple[str, Dict[str, Any]]], + max_changes: int, +) -> str: + now_utc = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S UTC") + shown = changes[:max_changes] + lines = [f"🚨 Wallet Activity ({len(changes)} changes):", ""] + + for idx, (change_type, pos) in enumerate(shown): + lines.append(_format_change_block(change_type, wallet, pos, now_utc)) + if idx != len(shown) - 1: + lines.append("") + + if len(changes) > max_changes: + lines.append("") + lines.append(f"... and {len(changes) - max_changes} more changes") + + return "\n".join(lines) + + +def start_polymarket_wallet_activity_loop(bot: Any) -> Optional[threading.Thread]: + enabled = _env_bool("POLYMARKET_WALLET_ACTIVITY_ENABLED", False) + chat_id = os.getenv("TELEGRAM_CHAT_ID") + users = _parse_addresses(os.getenv("POLYMARKET_WALLET_ACTIVITY_USERS")) + + if not enabled: + logger.info("polymarket wallet activity watcher disabled") + return None + if not chat_id: + logger.warning("polymarket wallet activity watcher skipped: TELEGRAM_CHAT_ID is not set") + return None + if not users: + logger.warning("polymarket wallet activity watcher skipped: POLYMARKET_WALLET_ACTIVITY_USERS is empty") + return None + + data_api_url = str( + os.getenv("POLYMARKET_WALLET_ACTIVITY_DATA_API_URL", "https://data-api.polymarket.com") + ).strip() + poll_sec = max(5, _env_int("POLYMARKET_WALLET_ACTIVITY_INTERVAL_SEC", 20)) + timeout_sec = max(5, _env_int("POLYMARKET_WALLET_ACTIVITY_TIMEOUT_SEC", 10)) + min_size_abs = max(0.0, _env_float("POLYMARKET_WALLET_ACTIVITY_MIN_SIZE_ABS", 0.001)) + min_size_delta = max(0.0, _env_float("POLYMARKET_WALLET_ACTIVITY_MIN_SIZE_DELTA", 0.001)) + max_changes = max(1, _env_int("POLYMARKET_WALLET_ACTIVITY_MAX_CHANGES_PER_MSG", 5)) + notify_closed = _env_bool("POLYMARKET_WALLET_ACTIVITY_NOTIFY_CLOSED", True) + bootstrap_alert = _env_bool("POLYMARKET_WALLET_ACTIVITY_BOOTSTRAP_ALERT", False) + + state_path = _state_file() + session = requests.Session() + + def _runner() -> None: + state = _load_state(state_path) + users_state = state.setdefault("users", {}) + + logger.info( + f"polymarket wallet activity watcher started users={len(users)} " + f"poll={poll_sec}s data_api={data_api_url} state_path={state_path}" + ) + + while True: + touched = False + for user in users: + try: + rows = _fetch_positions( + session=session, + base_url=data_api_url, + user=user, + timeout_sec=timeout_sec, + ) + current = _build_snapshot(rows, min_size_abs=min_size_abs) + prev = ( + (users_state.get(user) or {}).get("positions") + if isinstance(users_state.get(user), dict) + else {} + ) or {} + + if not prev and not bootstrap_alert: + users_state[user] = { + "positions": current, + "updated_at": int(time.time()), + } + touched = True + continue + + changes = _diff_positions( + previous=prev, + current=current, + min_size_delta=min_size_delta, + notify_closed=notify_closed, + ) + + if changes: + msg = _build_message(user, changes, max_changes=max_changes) + bot.send_message(chat_id, msg, disable_web_page_preview=True) + logger.info( + f"wallet activity pushed user={user} changes={len(changes)}" + ) + + users_state[user] = { + "positions": current, + "updated_at": int(time.time()), + } + touched = True + except Exception: + logger.exception(f"wallet activity cycle failed user={user}") + + if touched: + try: + _save_state(state_path, state) + except Exception: + logger.exception("failed to save wallet activity state") + + time.sleep(poll_sec) + + thread = threading.Thread( + target=_runner, + name="polymarket-wallet-activity-watcher", + daemon=True, + ) + thread.start() + return thread + +