Compare commits

..
Author SHA1 Message Date
schrodinger01andpselamy 96d1d3dd10 fix(alerter): escape all dynamic numeric fields in Telegram MarkdownV2
Telegram's MarkdownV2 parser treats `.` as a special character and rejects
the message with `Bad Request: can't parse entities` if any unescaped `.`
appears outside a code or pre block. The current formatter only escapes
the static `Market:` title and the `$` in the USDC amount, while leaving
three numeric runs unescaped:

  *Risk Score:* 0.82 (HIGH)            ← `.` in the score literal
  *Trade:* BUY Yes @ $0.075 | $15,000.00
                       ↑                  ↑
                       price              USDC amount

In production this means well-formed alerts silently fail to reach
Telegram users — the dispatcher reports a 400 from the Bot API and the
event is dropped.

This change routes every dynamic value through `_escape_telegram_markdown`
before interpolation:

  - `assessment.weighted_score` (risk score)
  - `trade.price` (price string)
  - `format_usdc(trade.notional_value)` (USDC amount, including its `.`)
  - `trade.side` and `trade.outcome` (defensive — upstream values may
    contain `-` or `.` in future schema changes)

The test that previously asserted `"0.82" in result.telegram_markdown` is
updated to require the escaped form `"0\.82"`, and a new test pins down
that none of `0.82` / `0.075` / `15,000.00` appear in unescaped form.
2026-06-14 19:10:34 +00:00
4 changed files with 53 additions and 422 deletions
@@ -290,12 +290,14 @@ class AlertFormatter:
wallet_line += f" \\(Age: {age_hours:.0f}h\\)"
lines.append(wallet_line)
# Risk score
lines.append(f"*Risk Score:* {assessment.weighted_score:.2f} \\({risk_level}\\)")
# Risk score — every numeric literal here must be escaped because
# MarkdownV2 treats `.` as a special character and rejects unescaped
# ones with `Bad Request: can't parse entities`.
score_str = self._escape_telegram_markdown(f"{assessment.weighted_score:.2f}")
lines.append(f"*Risk Score:* {score_str} \\({risk_level}\\)")
# Market
market_title = trade.event_title or trade.market_slug or "Unknown Market"
# Escape special Telegram markdown characters
market_title_escaped = self._escape_telegram_markdown(market_title)
if "market" in links:
lines.append(f"*Market:* [{market_title_escaped}]({links['market']})")
@@ -303,14 +305,16 @@ class AlertFormatter:
lines.append(f"*Market:* {market_title_escaped}")
# Trade details
usdc_value = format_usdc(trade.notional_value).replace("$", "\\$")
lines.append(
f"*Trade:* {trade.side} {trade.outcome} @ \\${trade.price:.3f} \\| {usdc_value}"
)
usdc_value = self._escape_telegram_markdown(format_usdc(trade.notional_value))
price_str = self._escape_telegram_markdown(f"{trade.price:.3f}")
side_escaped = self._escape_telegram_markdown(trade.side)
outcome_escaped = self._escape_telegram_markdown(trade.outcome)
lines.append(f"*Trade:* {side_escaped} {outcome_escaped} @ \\${price_str} \\| {usdc_value}")
# Signals
if signals:
lines.append(f"*Signals:* {', '.join(signals)}")
signals_escaped = [self._escape_telegram_markdown(s) for s in signals]
lines.append(f"*Signals:* {', '.join(signals_escaped)}")
# Links
lines.append("")
@@ -26,46 +26,8 @@ logger = logging.getLogger(__name__)
USDC_BRIDGED = "0x2791Bca1f2de4661ED88A30C99A7a9449Aa84174"
USDC_NATIVE = "0x3c499c542cEF5E3811e1192ce70d8cC03d5c3359"
# ERC20 Transfer event signature. ``HexBytes.hex()`` returns a *bare* hex
# string without the ``0x`` prefix; publicnode tolerates that, but stricter
# providers (e.g. drpc — which we use as the fallback) reject it with
# ``invalid argument 0: hex string without 0x prefix``. Always pass the
# 0x-prefixed form to ``eth_getLogs``.
# ERC20 Transfer event signature
TRANSFER_EVENT_SIGNATURE = AsyncWeb3.keccak(text="Transfer(address,address,uint256)")
TRANSFER_EVENT_TOPIC = "0x" + TRANSFER_EVENT_SIGNATURE.hex().removeprefix("0x")
# eth_getLogs block-range chunking. Most public Polygon RPCs (publicnode, ankr,
# llamarpc) cap the range at 10_000 blocks per call; pick a window slightly
# under the cap so off-by-one differences between providers don't trip us up.
DEFAULT_CHUNK_SIZE_BLOCKS = 9_000
# Polygon block time is ~2.0s. publicnode (the most common free RPC) prunes
# log history aggressively — empirically only ~100k blocks (~55 hours) are
# served before requests start returning "History has been pruned". We default
# to 80k blocks (~44 hours), which is more than enough for fresh-wallet
# funding traces (those wallets are by definition new) and fits comfortably
# inside what most public providers retain.
DEFAULT_MAX_LOOKBACK_BLOCKS = 80_000
# Substrings that, when present in an RPC error, indicate the chunk we just
# asked for is outside the provider's archive horizon. Walking further back
# is futile, so we stop the trace early instead of hammering every chunk.
_PRUNED_HISTORY_MARKERS: tuple[str, ...] = (
"history has been pruned",
"missing trie node",
"older than",
)
def _is_pruned_history_error(err: BaseException) -> bool:
"""Return True if the RPC error indicates pruned history.
Public Polygon nodes only retain a recent slice of log history. When we
walk back through that slice in chunks and hit the cutoff, every further
chunk will fail with the same message — so we stop early instead of
burning quota on guaranteed failures.
"""
text = str(err).lower()
return any(marker in text for marker in _PRUNED_HISTORY_MARKERS)
class FundingTracer:
@@ -88,8 +50,6 @@ class FundingTracer:
*,
max_hops: int = 3,
usdc_addresses: list[str] | None = None,
chunk_size_blocks: int = DEFAULT_CHUNK_SIZE_BLOCKS,
max_lookback_blocks: int = DEFAULT_MAX_LOOKBACK_BLOCKS,
) -> None:
"""Initialize the funding tracer.
@@ -98,11 +58,6 @@ class FundingTracer:
entity_registry: Registry for entity classification. Creates default if None.
max_hops: Maximum hops to trace back (default 3).
usdc_addresses: USDC contract addresses to track. Uses defaults if None.
chunk_size_blocks: Block window size per eth_getLogs call. Public
Polygon RPCs cap at 10_000 blocks; default leaves a safety margin.
max_lookback_blocks: How far back to scan when caller passes
``from_block=0``. Default ~44 hours at 2s block time, which
fits inside the pruned-history horizon of most public RPCs.
"""
self.polygon_client = polygon_client
self.entity_registry = entity_registry or EntityRegistry()
@@ -110,8 +65,6 @@ class FundingTracer:
self._usdc_addresses = [
addr.lower() for addr in (usdc_addresses or [USDC_BRIDGED, USDC_NATIVE])
]
self._chunk_size_blocks = chunk_size_blocks
self._max_lookback_blocks = max_lookback_blocks
async def trace(
self,
@@ -256,144 +209,46 @@ class FundingTracer:
) -> list[dict[str, Any]]:
"""Get ERC20 Transfer event logs.
Public Polygon RPCs (publicnode, ankr, llamarpc) cap ``eth_getLogs`` at
10_000 blocks per call. To work around this we resolve the requested
range into a concrete block window (defaulting to the last
``max_lookback_blocks`` when caller passes ``from_block=0``) and walk
the window in chunks of ``chunk_size_blocks``, oldest-first, stopping
once ``limit`` matches are collected. Walking oldest-first preserves
the "first transfer" semantics expected by the funding chain tracer.
If a chunk comes back with a "history has been pruned" style error
the rest of the walk is short-circuited — every subsequent chunk
would hit the same archive cutoff and there's no point burning quota
on guaranteed failures.
Args:
to_address: Filter by recipient address.
token_address: ERC20 token contract address.
limit: Maximum logs to return.
from_block: Starting block number (0 means
``latest - max_lookback_blocks``).
to_block: Ending block number ("latest" resolves to current head).
from_block: Starting block number.
to_block: Ending block number.
Returns:
List of log dictionaries, oldest first, capped at ``limit``.
List of log dictionaries.
"""
# Pad address to 32 bytes for topic filter
padded_to = "0x" + to_address.lower().replace("0x", "").zfill(64)
topics = [
TRANSFER_EVENT_TOPIC, # Transfer event (must be 0x-prefixed for drpc)
None, # from (any)
padded_to, # to (target address)
]
contract_address = AsyncWeb3.to_checksum_address(token_address)
start_block, end_block = await self._resolve_block_range(from_block, to_block)
if start_block > end_block:
return []
results: list[dict[str, Any]] = []
chunk_size = max(1, self._chunk_size_blocks)
chunk_start = start_block
while chunk_start <= end_block:
chunk_end = min(chunk_start + chunk_size - 1, end_block)
try:
chunk_logs = await self._fetch_logs_chunk(
contract_address=contract_address,
topics=topics,
from_block=chunk_start,
to_block=chunk_end,
)
except Exception as e:
if _is_pruned_history_error(e):
# The provider has dropped this slice of history. Walking
# further back will hit the same wall on every chunk;
# stop now and return what we already have.
logger.info(
"eth_getLogs chunk %d-%d outside archive horizon for %s; stopping trace",
chunk_start,
chunk_end,
to_address,
)
break
logger.warning(
"eth_getLogs chunk %d-%d failed for %s: %s",
chunk_start,
chunk_end,
to_address,
e,
)
# Skip this window and keep walking — partial data is better
# than aborting the whole trace on a single flaky chunk.
chunk_start = chunk_end + 1
continue
for log in chunk_logs:
results.append(dict(log))
if len(results) >= limit:
return results
chunk_start = chunk_end + 1
return results
async def _resolve_block_range(
self,
from_block: int | str,
to_block: int | str,
) -> tuple[int, int]:
"""Resolve symbolic block params to concrete numeric bounds.
``from_block=0`` (the historical default) is rewritten to
``latest - max_lookback_blocks`` so we don't try to scan all of Polygon.
"""
w3 = self._select_w3()
if isinstance(to_block, str):
await self.polygon_client._rate_limiter.acquire()
head = int(await w3.eth.block_number)
end = head
else:
end = int(to_block)
if isinstance(from_block, str):
# Treat any symbolic from-block (e.g. "earliest") as "go back
# max_lookback_blocks from end"; that's what callers actually want.
start = max(0, end - self._max_lookback_blocks)
elif from_block == 0:
start = max(0, end - self._max_lookback_blocks)
else:
start = int(from_block)
return start, end
async def _fetch_logs_chunk(
self,
contract_address: str,
topics: list[Any],
from_block: int,
to_block: int,
) -> list[Any]:
"""Issue a single bounded ``eth_getLogs`` call."""
await self.polygon_client._rate_limiter.acquire()
w3 = self._select_w3()
# Use the web3 instance from polygon client
w3 = (
self.polygon_client._w3
if self.polygon_client._primary_healthy
else (self.polygon_client._w3_fallback or self.polygon_client._w3)
)
# Get logs with Transfer event filtering by recipient
# Note: web3 typing is overly restrictive for block params
return await w3.eth.get_logs(
logs = await w3.eth.get_logs(
{
"address": contract_address,
"topics": topics,
"address": AsyncWeb3.to_checksum_address(token_address),
"topics": [
TRANSFER_EVENT_SIGNATURE.hex(), # Transfer event
None, # from (any)
padded_to, # to (target address)
],
"fromBlock": from_block, # type: ignore[typeddict-item]
"toBlock": to_block, # type: ignore[typeddict-item]
}
)
def _select_w3(self) -> AsyncWeb3:
"""Pick primary or fallback web3 instance based on health."""
if self.polygon_client._primary_healthy:
return self.polygon_client._w3
return self.polygon_client._w3_fallback or self.polygon_client._w3
# Convert to list of dicts and limit
result = [dict(log) for log in logs[:limit]]
return result
async def _log_to_funding_transfer(
self,
+17 -2
View File
@@ -409,12 +409,27 @@ class TestTelegramMarkdown:
assert "`0x1234...5678`" in result.telegram_markdown
def test_telegram_includes_risk_score(self, high_risk_assessment: RiskAssessment) -> None:
"""Test that Telegram message includes risk score."""
"""Test that Telegram message includes risk score, MarkdownV2-escaped."""
formatter = AlertFormatter()
result = formatter.format(high_risk_assessment)
assert "0.82" in result.telegram_markdown
# MarkdownV2 requires `.` to be escaped, so 0.82 becomes 0\.82.
assert "0\\.82" in result.telegram_markdown
assert "HIGH" in result.telegram_markdown
def test_telegram_escapes_all_decimals(self, high_risk_assessment: RiskAssessment) -> None:
"""Telegram MarkdownV2 rejects unescaped `.` in dynamic numeric
fields (risk score, price, USDC amount). All must be present in
their escaped form."""
formatter = AlertFormatter()
result = formatter.format(high_risk_assessment)
md = result.telegram_markdown
for unescaped in ("0.82", "0.075", "15,000.00"):
assert unescaped not in md, (
f"unescaped {unescaped!r} would be rejected by Telegram MarkdownV2: {md!r}"
)
for escaped in ("0\\.82", "0\\.075", "15,000\\.00"):
assert escaped in md, f"missing escaped {escaped!r} in {md!r}"
def test_telegram_includes_links(self, high_risk_assessment: RiskAssessment) -> None:
"""Test that Telegram message includes links."""
formatter = AlertFormatter()
+1 -244
View File
@@ -327,10 +327,6 @@ class TestGetTransferLogs:
await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
# Explicit numeric range so we stay inside one chunk and skip
# the "latest" → block_number resolution path.
from_block=1,
to_block=8_000,
)
mock_w3.eth.get_logs.assert_called_once()
@@ -338,15 +334,10 @@ class TestGetTransferLogs:
# Verify topics structure
assert len(call_args["topics"]) == 3
# The Transfer event topic must be 0x-prefixed; drpc rejects bare hex.
assert call_args["topics"][0] == "0x" + TRANSFER_EVENT_SIGNATURE.hex().removeprefix("0x")
assert call_args["topics"][0].startswith("0x")
assert call_args["topics"][0] == TRANSFER_EVENT_SIGNATURE.hex()
assert call_args["topics"][1] is None # from (any)
# to address should be padded to 32 bytes
assert call_args["topics"][2].endswith(TEST_WALLET.lower().replace("0x", ""))
# And the chunk bounds match what we asked for.
assert call_args["fromBlock"] == 1
assert call_args["toBlock"] == 8_000
@pytest.mark.asyncio
async def test_get_transfer_logs_respects_limit(
@@ -364,8 +355,6 @@ class TestGetTransferLogs:
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
limit=3,
from_block=1,
to_block=8_000,
)
assert len(result) == 3
@@ -385,242 +374,10 @@ class TestGetTransferLogs:
await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
from_block=1,
to_block=8_000,
)
mock_fallback.eth.get_logs.assert_called_once()
@pytest.mark.asyncio
async def test_get_transfer_logs_chunks_large_ranges(
self,
funding_tracer: FundingTracer,
mock_polygon_client: MagicMock,
) -> None:
"""Ranges wider than chunk_size are split into multiple eth_getLogs calls.
This is the regression guard for the publicnode 10_000-block cap that
was rejecting every funding trace before chunking landed.
"""
mock_w3 = MagicMock()
mock_w3.eth.get_logs = AsyncMock(return_value=[])
mock_polygon_client._w3 = mock_w3
# 25_000 blocks at 9_000-per-chunk → 3 calls (9000 + 9000 + 7001).
await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
from_block=1_000_000,
to_block=1_025_000,
)
assert mock_w3.eth.get_logs.call_count == 3
windows = [call[0][0] for call in mock_w3.eth.get_logs.call_args_list]
assert windows[0]["fromBlock"] == 1_000_000
assert windows[0]["toBlock"] == 1_008_999
assert windows[1]["fromBlock"] == 1_009_000
assert windows[1]["toBlock"] == 1_017_999
assert windows[2]["fromBlock"] == 1_018_000
assert windows[2]["toBlock"] == 1_025_000
# No window exceeds the chunk size — that's what RPC providers reject.
for win in windows:
assert win["toBlock"] - win["fromBlock"] + 1 <= 9_000
@pytest.mark.asyncio
async def test_get_transfer_logs_stops_when_limit_hit_mid_walk(
self,
funding_tracer: FundingTracer,
mock_polygon_client: MagicMock,
) -> None:
"""Walking should stop as soon as ``limit`` matches are gathered."""
mock_w3 = MagicMock()
# First chunk yields 5 logs, more than the limit, so subsequent chunks
# must not be queried.
mock_w3.eth.get_logs = AsyncMock(return_value=[MagicMock() for _ in range(5)])
mock_polygon_client._w3 = mock_w3
result = await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
limit=2,
from_block=1_000_000,
to_block=1_025_000,
)
assert len(result) == 2
mock_w3.eth.get_logs.assert_called_once()
@pytest.mark.asyncio
async def test_get_transfer_logs_skips_failing_chunk(
self,
funding_tracer: FundingTracer,
mock_polygon_client: MagicMock,
) -> None:
"""A flaky chunk must not abort the whole trace — we move on."""
mock_w3 = MagicMock()
good_log = MagicMock()
responses: list[Any] = [
RuntimeError("RPC hiccup"),
[good_log],
]
async def fake_get_logs(_params: dict[str, Any]) -> list[Any]:
outcome = responses.pop(0)
if isinstance(outcome, BaseException):
raise outcome
return outcome
mock_w3.eth.get_logs = AsyncMock(side_effect=fake_get_logs)
mock_polygon_client._w3 = mock_w3
result = await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
from_block=1_000_000,
to_block=1_018_000, # forces 3 chunks; we exercise chunks 1+2
)
# The error chunk is skipped; the second chunk contributes one log.
assert result == [dict(good_log)]
assert mock_w3.eth.get_logs.call_count >= 2
@pytest.mark.asyncio
async def test_get_transfer_logs_resolves_latest_via_block_number(
self,
funding_tracer: FundingTracer,
mock_polygon_client: MagicMock,
) -> None:
"""``to_block='latest'`` should resolve via ``eth.block_number``.
And ``from_block=0`` should not become a full-history scan — it must
be clamped to ``latest - max_lookback_blocks``.
"""
async def _block_number_coro() -> int:
return 5_000
mock_eth = MagicMock()
mock_eth.get_logs = AsyncMock(return_value=[])
# Property-style awaitable: web3.py exposes block_number as a property
# returning a coroutine, so each access must yield a fresh awaitable.
type(mock_eth).block_number = property( # type: ignore[misc]
lambda _self: _block_number_coro()
)
mock_w3 = MagicMock()
mock_w3.eth = mock_eth
mock_polygon_client._w3 = mock_w3
await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
)
# block_number=5000 < chunk_size, so it's one chunk that bottoms at 0.
mock_eth.get_logs.assert_called_once()
call_args = mock_eth.get_logs.call_args[0][0]
assert call_args["fromBlock"] == 0
assert call_args["toBlock"] == 5_000
@pytest.mark.asyncio
async def test_get_transfer_logs_breaks_on_pruned_history(
self,
funding_tracer: FundingTracer,
mock_polygon_client: MagicMock,
) -> None:
"""A pruned-history error must short-circuit the whole walk.
Public Polygon RPCs prune log history. Once we walk past the cutoff,
every subsequent chunk will raise the same error — keep walking and
we just burn quota on guaranteed failures. The first such error must
end the walk and return whatever we already collected.
"""
mock_w3 = MagicMock()
good_log = MagicMock()
responses: list[Any] = [
[good_log],
RuntimeError(
"{'code': -32701, 'message': 'History has been pruned for "
"this block. To remove restrictions, order a dedicated full "
"node here: https://www.allnodes.com/pol/host'}"
),
# If the early-break logic is missing, this third chunk would
# also be requested. The test asserts it isn't.
[MagicMock()],
]
async def fake_get_logs(_params: dict[str, Any]) -> list[Any]:
outcome = responses.pop(0)
if isinstance(outcome, BaseException):
raise outcome
return outcome
mock_w3.eth.get_logs = AsyncMock(side_effect=fake_get_logs)
mock_polygon_client._w3 = mock_w3
# 3 chunks total. The pruned error fires on chunk #2; chunk #3 must
# never be issued.
result = await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
from_block=1_000_000,
to_block=1_027_000,
)
assert result == [dict(good_log)]
assert mock_w3.eth.get_logs.call_count == 2
@pytest.mark.asyncio
async def test_get_transfer_logs_default_lookback_fits_pruned_horizon(
self,
) -> None:
"""Default ``max_lookback_blocks`` must stay inside what public RPCs serve.
publicnode prunes after ~100k blocks. If we default to 1.3M, every
funding trace blows through the archive horizon and produces nothing
but pruned-history warnings. Pin the default at <= 100k as a
regression guard.
"""
from polymarket_insider_tracker.profiler.funding import (
DEFAULT_MAX_LOOKBACK_BLOCKS,
)
assert DEFAULT_MAX_LOOKBACK_BLOCKS <= 100_000
@pytest.mark.asyncio
async def test_get_transfer_logs_topic_is_0x_prefixed(
self,
funding_tracer: FundingTracer,
mock_polygon_client: MagicMock,
) -> None:
"""The Transfer event topic passed to ``eth_getLogs`` must begin with ``0x``.
``HexBytes.hex()`` returns a bare hex string. publicnode tolerates
that, but stricter providers like drpc (our fallback) reject it with
``invalid argument 0: hex string without 0x prefix`` and every chunk
in the trace fails. This guards against regressing back to the
bare-hex form.
"""
mock_w3 = MagicMock()
mock_w3.eth.get_logs = AsyncMock(return_value=[])
mock_polygon_client._w3 = mock_w3
await funding_tracer._get_transfer_logs(
to_address=TEST_WALLET,
token_address=USDC_BRIDGED,
from_block=1,
to_block=8_000,
)
topics = mock_w3.eth.get_logs.call_args[0][0]["topics"]
assert topics[0].startswith("0x")
# And the topic also has to be 32 bytes (64 hex chars) as required by
# the JSON-RPC spec.
assert len(topics[0]) == 2 + 64
# The padded `to` topic was already 0x-prefixed; double-check that
# didn't regress either.
assert topics[2].startswith("0x")
class TestLogToFundingTransfer:
"""Tests for _log_to_funding_transfer method."""