Garden-Angel-Ai-35Bot / modules /feed_client.py
35
two /doctor lies that only appear once Redis is up (#178)
7d2e91e unverified
Raw
History Blame Contribute Delete
24.5 kB
"""
modules/feed_client.py — read the shared Redis pool feed. Trigger, not price.
────────────────────────────────────────────────────────────────────────────
WHAT PROBLEM THIS SOLVES
────────────────────────────────────────────────────────────────────────────
/doctor reports 6,927/7,000 Jupiter calls used per hour (99%) on four
routes, and /capture shows per-route pipeline times of 959ms–2,346ms while
/pipeline's own p50 is 100ms. That gap is the scanner waiting for quote
budget it has already spent.
Almost all of that spend is asking Jupiter a question whose answer has not
changed. The bot polls every route on a timer; the pools underneath it move
only when someone trades them. A separate service (`garden-feed`, written by
the Rust developer, NOT started or managed from here) streams pool state
from Triton Dragon's Mouth into Redis. This module reads it.
────────────────────────────────────────────────────────────────────────────
TRIGGER, NOT PRICE — AND WHY THAT DISTINCTION IS THE WHOLE DESIGN
────────────────────────────────────────────────────────────────────────────
The obvious use of a reserves feed is to compute the swap output locally and
skip Jupiter entirely. For Raydium's classic constant-product pools that is
correct and easy: x·y=k with a 25 bps fee, and reserve_a/reserve_b are
genuinely sufficient.
It is WRONG for every other venue this bot trades. /venues shows the real
routes going through Whirlpool, Meteora DLMM and Raydium CLMM — all
concentrated liquidity. In a CLMM the price comes from the current tick and
the shape of liquidity around it, not from total reserves. Two pools holding
identical reserve_a and reserve_b can quote completely different prices
depending on where that liquidity sits. There is no formula that turns
{reserve_a, reserve_b} into a correct CLMM quote, and a number computed that
way would be confidently, silently wrong on the majority of this
deployment's volume — on real money.
So this module never computes a price. It answers exactly one question:
HAS THIS POOL MOVED SINCE THE LAST TIME I ASKED JUPITER ABOUT IT?
If nothing moved, the previous answer still stands and the quote is skipped.
If something moved, Jupiter is asked — and Jupiter's answer is the price, as
it is today. The saving comes from not asking when nothing changed, which is
where the budget was actually going. It is correct on every DEX from the
first cycle, and it cannot lose money by being wrong about AMM maths,
because it does no AMM maths.
(If the feed later publishes sqrt_price / liquidity / tick_current, local
CLMM pricing becomes possible and this module is where it would go. That is
a change to the feed contract and belongs to the operator, not to this file.)
────────────────────────────────────────────────────────────────────────────
CONTRACT — matches the feed-writer's brief exactly
────────────────────────────────────────────────────────────────────────────
feed:sol-usdc:<dex_slug>:state JSON {reserve_a, reserve_b, mint_a,
mint_b, slot, ts_ms}
feed:sol-usdc:pools SET of live dex_slugs
`ts_ms` is the writer's local epoch-millis, not chain time. Freshness is
judged from it here, as the writer's brief specifies.
Two briefs were issued with DIFFERENT key shapes — the writer's
`feed:sol-usdc:<dex>:state` and a consumer draft's `feed:solusdc:<dex>`.
This follows the WRITER, because a consumer that reads a key nobody writes
finds nothing, falls back to Jupiter forever, and looks perfectly healthy
while doing it. FEED_KEY_PREFIX exists so that can be corrected without a
deploy if the contract moves.
────────────────────────────────────────────────────────────────────────────
RULES THIS MODULE WILL NOT BREAK
────────────────────────────────────────────────────────────────────────────
IT NEVER CRASHES THE BOT. Redis unreachable, malformed JSON, missing key,
wrong types — every one of them degrades to "I don't know", and "I don't
know" always means "ask Jupiter, exactly as before". A pricing optimisation
that can stop a bot trading is not one.
IT DOES NOT MANAGE THE FEED. It never starts, restarts, or health-checks
`garden-feed`. Read-only consumer.
IT IS OFF BY DEFAULT. FEED_ENABLED=false, so old and new behaviour can be
compared without a redeploy, and an upgrade changes nothing until asked.
FEED_ENABLED=true turn it on
FEED_MAX_AGE_MS=250 older than this and the value is not trusted
FEED_STALE_ACTION=skip skip | jupiter (default skip, per brief)
FEED_REDIS_URL=redis://127.0.0.1:6379/0
FEED_KEY_PREFIX=feed:sol-usdc
"""
from __future__ import annotations
import json
import logging
import os
import time
from dataclasses import dataclass
from typing import Any, Optional
logger = logging.getLogger(__name__)
def _env_bool(name: str, default: bool) -> bool:
raw = os.getenv(name, "").strip().lower()
return raw in ("1", "true", "yes", "on") if raw else default
def _env_float(name: str, default: float) -> float:
try:
return float(os.getenv(name, "").strip() or default)
except (TypeError, ValueError):
return default
ENABLED = _env_bool("FEED_ENABLED", False)
MAX_AGE_MS = _env_float("FEED_MAX_AGE_MS", 250.0)
STALE_ACTION = (os.getenv("FEED_STALE_ACTION", "").strip().lower() or "skip")
# When /doctor should call the writer stopped rather than the market quiet.
# Deliberately far above MAX_AGE_MS: that one gates a single trade and 250ms
# is right for it, but a diagnostic that shouts at 250ms would shout at every
# lull. Thirty seconds of nothing from EVERY pool is not a lull.
_STALL_AFTER_MS = max(30_000.0, MAX_AGE_MS * 30.0)
REDIS_URL = os.getenv("FEED_REDIS_URL", "").strip() or "redis://127.0.0.1:6379/0"
KEY_PREFIX = os.getenv("FEED_KEY_PREFIX", "").strip() or "feed:sol-usdc"
# How long to stay quiet after a Redis failure before trying again. Without
# this, an unreachable Redis turns into a connection attempt on every route
# of every cycle — replacing the latency this module exists to remove.
_RETRY_AFTER_SECS = _env_float("FEED_RETRY_AFTER_SECS", 30.0)
@dataclass
class PoolState:
dex: str
reserve_a: int
reserve_b: int
slot: int
ts_ms: int
mint_a: str = ""
mint_b: str = ""
def age_ms(self) -> float:
return max(0.0, time.time() * 1000.0 - self.ts_ms)
def fresh(self) -> bool:
return self.age_ms() <= MAX_AGE_MS
def fingerprint(self) -> tuple:
"""What 'this pool moved' means, concretely.
Reserves AND slot. Reserves alone would miss a trade that happens to
leave both sides numerically identical (rare, but a rounding-sized
swap can do it); slot alone changes constantly whether or not this
pool was touched, and would defeat the entire point by making every
pool look moved on every slot.
"""
return (self.reserve_a, self.reserve_b, self.slot)
class FeedClient:
"""Read-only view of the shared pool feed. Never raises outward."""
def __init__(self) -> None:
self._redis: Any = None
self._unavailable_until = 0.0
self._last_error = ""
# dex -> the fingerprint at the moment Jupiter was last asked.
self._seen: dict[str, tuple] = {}
self._hits = 0
self._skips = 0
self._stale = 0
self._misses = 0
# ── connection ───────────────────────────────────────────────────────
def _client(self) -> Any:
if not ENABLED:
return None
if self._redis is not None:
return self._redis
if time.monotonic() < self._unavailable_until:
return None # still in the quiet period
try:
import redis # type: ignore
client = redis.Redis.from_url(
REDIS_URL, socket_timeout=0.25, socket_connect_timeout=0.25,
decode_responses=True,
)
client.ping()
self._redis = client
self._last_error = ""
logger.info("[Feed] connected to %s — pool feed is live", REDIS_URL)
return client
except Exception as exc: # noqa: BLE001 — never fatal, see docstring
self._note_unavailable(exc)
return None
def _note_unavailable(self, exc: BaseException) -> None:
self._redis = None
first = not self._last_error
self._last_error = str(exc)[:200]
self._unavailable_until = time.monotonic() + _RETRY_AFTER_SECS
# Logged once per outage, not once per lookup: a down Redis on a
# 4-route cycle would otherwise produce thousands of identical lines
# an hour and bury everything else.
if first:
logger.warning(
"[Feed] Redis unavailable (%s) — falling back to the normal "
"Jupiter path. Retrying in %.0fs. Trading is unaffected.",
self._last_error, _RETRY_AFTER_SECS,
)
# ── reads ────────────────────────────────────────────────────────────
def pools(self) -> list[str]:
"""Live dex slugs, straight from the writer. Never hardcoded here."""
client = self._client()
if client is None:
return []
try:
return sorted(str(m) for m in client.smembers(f"{KEY_PREFIX}:pools"))
except Exception as exc: # noqa: BLE001
self._note_unavailable(exc)
return []
def state(self, dex: str) -> Optional[PoolState]:
client = self._client()
if client is None:
return None
try:
raw = client.get(f"{KEY_PREFIX}:{dex}:state")
except Exception as exc: # noqa: BLE001
self._note_unavailable(exc)
return None
if not raw:
self._misses += 1
return None
try:
d = json.loads(raw)
return PoolState(
dex=dex,
reserve_a=int(d["reserve_a"]), reserve_b=int(d["reserve_b"]),
slot=int(d.get("slot") or 0), ts_ms=int(d["ts_ms"]),
mint_a=str(d.get("mint_a") or ""), mint_b=str(d.get("mint_b") or ""),
)
except Exception as exc: # noqa: BLE001 — a bad row is a miss, not a crash
self._misses += 1
logger.debug("[Feed] unparsable state for %s: %s", dex, exc)
return None
def describe(self) -> dict[str, Any]:
"""What is ACTUALLY in Redis, versus what we are reading. For /doctor.
────────────────────────────────────────────────────────────────────
THE FAILURE THIS EXISTS TO MAKE LOUD
────────────────────────────────────────────────────────────────────
Two briefs specified two different key layouts for the same data:
feed-writer (Rust) : feed:sol-usdc:<dex>:state
consumer draft : feed:solusdc:<dex>
`sol-usdc` vs `solusdc`, and `:state` present vs absent. This module
follows the WRITER's version, because the writer is the side that
decides what exists. But if that reconciliation never happens — or
the Rust side changes it later — the consumer reads a key nobody
writes, finds nothing, and correctly falls back to Jupiter.
Correctly, and FOREVER, and silently. Every screen would report a
healthy bot with a healthy feed configured, the Jupiter budget would
stay at 99%, and the only symptom would be an improvement that never
arrived. There is no error to see, because nothing errors: a feed
with no data is indistinguishable from a quiet market, and the
safety asymmetry in should_quote() guarantees it stays quiet.
So this goes and LOOKS. It SCANs a bounded number of `feed:*` keys
and compares what is there to what we are configured to read. If the
writer is publishing under a different prefix, that prefix is on
screen with the one-line command that fixes it — no deploy, no code
change, and no waiting for two developers to compare documents.
Bounded and read-only: one SCAN with COUNT, capped at 200 keys, no
KEYS command (which blocks the server), and it never runs on the
trading path — only when a human asks.
"""
# v1.88 — PROBLEMS vs WARNINGS, and this module got it wrong.
#
# /doctor's ⛔ list is for things that stop the bot trading. It ran
# this on the live box and produced:
#
# ⛔ 4 problem(s):
# • signing key unavailable ... IAM, NOT SECRETS ...
# • pool feed: Redis unreachable ... trading is unaffected ...
#
# The second entry says "trading is unaffected" IN ITS OWN TEXT and
# was still filed as a blocker, sitting directly above the one thing
# that genuinely is one. The operator's instruction is "dont want
# see any warn"; a false blocker is worse than a warning, because it
# dilutes the list you are meant to read top-down.
#
# An unreachable Redis is EXPECTED here: the feed writer does not
# exist yet (the Rust side confirmed it), so nothing is lost. The
# only genuine problem this module can detect is the KEY MISMATCH
# below — data being published that we are not reading — because
# that one is silently costing something.
out: dict[str, Any] = {
"enabled": ENABLED, "configured_prefix": KEY_PREFIX,
"reachable": False, "pools": [], "found_prefixes": [],
"problems": [], "warnings": [], "info": "", "hint": "",
}
if not ENABLED:
out["info"] = "pool feed off (FEED_ENABLED=false) — Jupiter only"
return out
client = self._client()
if client is None:
out["warnings"].append(
f"pool feed is ON but Redis is unreachable at {REDIS_URL} "
f"({self._last_error}). Not blocking — the bot is on the "
f"normal Jupiter path, which is where it has always been. "
f"Expected until garden-feed exists; "
f"`./scripts/setenv.sh FEED_ENABLED=false` silences it."
)
return out
out["reachable"] = True
out["pools"] = self.pools()
try:
seen: set[str] = set()
cursor, scanned = 0, 0
while scanned < 200:
cursor, batch = client.scan(cursor=cursor, match="feed:*", count=100)
for key in batch:
scanned += 1
bits = str(key).split(":")
# feed:<pair>:<dex>:state -> prefix is the first two.
if len(bits) >= 2:
seen.add(":".join(bits[:2]))
if cursor == 0:
break
out["found_prefixes"] = sorted(seen)
except Exception as exc: # noqa: BLE001 — a diagnostic must not raise
out["warnings"].append(
f"pool feed: could not scan Redis: {str(exc)[:100]}")
return out
if out["pools"]:
# v1.89 — A REGISTERED POOL IS NOT A LIVE ONE.
#
# This returned "live" on the strength of the `:pools` set
# alone, and that set is written once at writer startup. So a
# writer that registered two pools and then died reported
# `pool feed live — 2 pool(s)` forever, with no warning and no
# problem. Verified against a real Redis: seed two states, wait
# one second (4x FEED_MAX_AGE_MS), ask /doctor — still "live".
#
# Same shape as the key mismatch this module was written to
# catch, one step further in: the feed looks healthy and is
# feeding nothing, and should_quote()'s staleness check sends
# everything to Jupiter exactly as it would with no feed at all.
# The improvement simply never arrives.
#
# The test is the FRESHEST pool, not each pool. One pool going
# quiet is a quiet market — a pool with no swaps has no update
# to publish, and warning about that would be a warning about
# normal. Every pool stale at once is a stopped writer.
newest_ms = None
for dex in out["pools"]:
st = self.state(dex)
ts = getattr(st, "ts_ms", None) if st else None
if ts and (newest_ms is None or ts > newest_ms):
newest_ms = ts
age_ms = ((time.time() * 1000.0) - newest_ms
if newest_ms else None)
if age_ms is not None and age_ms > _STALL_AFTER_MS:
out["warnings"].append(
f"pool feed: {len(out['pools'])} pool(s) are registered "
f"under `{KEY_PREFIX}` but the freshest one was written "
f"{age_ms / 1000.0:,.0f}s ago — the writer has stopped "
f"publishing. Not blocking: every route falls back to "
f"Jupiter, which is where it was before the feed. "
f"Check garden-feed is running."
)
elif newest_ms is None:
out["warnings"].append(
f"pool feed: {len(out['pools'])} pool(s) are registered "
f"under `{KEY_PREFIX}` but not one has a readable state "
f"key. The writer registered its pools and published no "
f"data. Not blocking; Jupiter as usual."
)
else:
out["info"] = (
f"pool feed live — {len(out['pools'])} pool(s) under "
f"`{KEY_PREFIX}`, freshest {age_ms:,.0f}ms old")
return out
if not out["found_prefixes"]:
# Redis up, writer absent. Still not a blocker: this is the
# documented state until garden-feed exists, and the fallback
# to Jupiter is the behaviour the bot has always had.
out["warnings"].append(
"pool feed: it is ON, Redis is up, and there is no feed data in "
"it at all — the writer is not running. Not blocking; the "
"bot quotes through Jupiter as usual."
)
else:
others = [p for p in out["found_prefixes"] if p != KEY_PREFIX]
if others:
out["problems"].append(
f"KEY MISMATCH — this bot reads `{KEY_PREFIX}:*` and finds "
f"nothing, but the writer is publishing under "
f"`{others[0]}:*`. The feed looks healthy and is buying "
f"nothing. This is the silent failure the two briefs "
f"disagreed about."
)
out["hint"] = f"./scripts/setenv.sh FEED_KEY_PREFIX={others[0]}"
else:
# NOT a blocker, and calling it one contradicted this
# module's own contract two screens above. A writer that has
# started and not yet registered a pool is a writer that is
# starting — the correct response is to wait, not to act.
out["warnings"].append(
f"pool feed: the prefix `{KEY_PREFIX}` exists but its "
f"`:pools` set is empty — the writer has started and has "
f"not registered any pool yet. Not blocking."
)
return out
# ── the only question this module answers ────────────────────────────
def should_quote(self, route_key: str) -> tuple[bool, str]:
"""(ask Jupiter?, why). Defaults to YES on every uncertainty.
The asymmetry is deliberate and is the safety property of this whole
module: a wrong "yes" costs one quote, a wrong "no" costs a trade.
Every unknown — feed off, Redis down, pool absent, JSON malformed,
value stale under FEED_STALE_ACTION=jupiter — resolves to yes.
The only "no" comes from positive evidence: the feed answered, the
value is fresh, and its fingerprint is byte-identical to the one
present when Jupiter was last asked about this route.
"""
if not ENABLED:
return True, "feed disabled"
pools = self.pools()
if not pools:
return True, "feed has no pools (or Redis is down)"
states = [s for s in (self.state(p) for p in pools) if s is not None]
if not states:
self._misses += 1
return True, "no pool state readable"
stale = [s for s in states if not s.fresh()]
if stale:
self._stale += 1
oldest = max(s.age_ms() for s in stale)
if STALE_ACTION == "jupiter":
return True, (f"feed stale ({oldest:.0f}ms > {MAX_AGE_MS:.0f}ms) "
f"— falling back to Jupiter")
# Default: skip. A stale feed means the writer is behind, and a
# price derived from behind-the-market data is worse than no
# trade. Per the operator's brief, skip is the default.
self._skips += 1
return False, (f"feed stale ({oldest:.0f}ms > {MAX_AGE_MS:.0f}ms) "
f"— skipping this cycle (FEED_STALE_ACTION=skip)")
fingerprint = tuple(sorted((s.dex, *s.fingerprint()) for s in states))
if self._seen.get(route_key) == fingerprint:
self._skips += 1
return False, "no pool moved since the last quote"
self._seen[route_key] = fingerprint
self._hits += 1
return True, "pool state changed"
# ── introspection for /doctor and /verify ────────────────────────────
def status(self) -> dict[str, Any]:
asked = self._hits + self._skips
return {
"enabled": ENABLED,
"connected": self._redis is not None,
"url": REDIS_URL,
"max_age_ms": MAX_AGE_MS,
"stale_action": STALE_ACTION,
"quotes_triggered": self._hits,
"quotes_skipped": self._skips,
"skip_rate": round(self._skips / asked, 3) if asked else None,
"stale_reads": self._stale,
"misses": self._misses,
"last_error": self._last_error,
}
_feed: Optional[FeedClient] = None
def get_feed() -> FeedClient:
global _feed
if _feed is None:
_feed = FeedClient()
return _feed