Garden-Angel-Ai-35Bot / modules /decay_tracker.py
35
Break the decay forever-loop, and stop /why quoting a 4-day-old cause (#153)
a91f0af unverified
Raw
History Blame Contribute Delete
28.8 kB
"""
modules/decay_tracker.py — measure how fast an edge dies between detection
and send, and gate on it (v1.0, 2026-07-29).
Operator, after 1,093 scan cycles and 0 signals: "Where's the signal we are
at Frankfurt zone". And, correctly, from their own analysis:
"Log (quote_profit, requote_profit, elapsed_ms) on every attempt, hit or
miss — right now you're tuning blind. That dataset tells you real
decay, not guessed decay."
Catch condition:
Net_at_detection ≥ MIN_PROFIT_FLOOR + Decay_Rate × T_pipeline
"The lever that actually catches more chances is shrinking T_pipeline,
not lowering MIN_PROFIT_FLOOR. Lowering the floor just admits more
trades that this same equation says you'll lose to decay anyway."
That is exactly right, and it is the piece this codebase never had. Worth
being precise about why it matters here.
────────────────────────────────────────────────────────────────────────────
THE PATTERN IN THIS DEPLOYMENT'S OWN HISTORY
────────────────────────────────────────────────────────────────────────────
Read solana_arb.py's changelog and the same event recurs under four
different diagnoses:
v1.42 a real +$1.51 signal's safety re-quote queued behind the rest of a
16-route cycle; ~20s later the edge was -$1.04 and it aborted.
v1.47 the scanner quoted leg 2 for the full leg-1 output while the
executor quoted it minus a haircut — signals that could never
reproduce at send time.
v1.16 a genuine +1.0 bps / +$1.95 signal died to a 429 moments after the
scan loop tripped the shared cooldown.
v1.50 a re-quote netting +$0.01 passed the margin gate and landed as a
real loss once network costs were counted.
Every one was found by hand, from logs, after the fact. Every one is the
same shape: **the number that justified the trade was measured at time T,
and the trade happened at time T + Δ.** Nothing in the bot has ever
measured Δ, or what it costs.
So the floor was tuned blind. $10 → $0.20 → $0.03 across this changelog,
each move a guess at "let more through", when the binding constraint was
never the height of the bar — it was that the edge was already gone by the
time anything jumped.
────────────────────────────────────────────────────────────────────────────
WHAT THIS MODULE DOES
────────────────────────────────────────────────────────────────────────────
1. RECORDS every real execution attempt as one sample:
quote_net_usd what the scanner saw, at detection
requote_net_usd what the mandatory pre-send re-quote saw
elapsed_ms detection → re-quote, wall clock
route, outcome
These three numbers together are the dataset that has never existed.
They are appended to their own CSV (separate from trade_journal.csv,
which is one row per EVALUATION — this is one row per real ATTEMPT, a
thousand times rarer and a different unit of analysis).
2. ESTIMATES the decay rate in USD per second, per route and globally, as
the median of (quote_net − requote_net) / elapsed_secs over recent
samples. Median, not mean: a single 429-delayed attempt with a 40s
elapsed would otherwise dominate the estimate for hours.
3. GATES. `required_floor()` returns
MIN_PROFIT_FLOOR + decay_rate × T_pipeline_p90
so a signal must clear not just the floor, but the floor plus what the
edge is expected to lose while the pipeline runs. A signal that cannot
clear it was never going to land profitably — attempting it burns
Jupiter budget and, at worst, real fees.
────────────────────────────────────────────────────────────────────────────
THE HONEST LIMITS
────────────────────────────────────────────────────────────────────────────
Decay is not one number. It is a mixture of at least three processes: other
searchers taking the same edge (fast, adversarial), ordinary price drift
(slow, symmetric), and our own pipeline stalling (bursty, self-inflicted).
A median lumps them together. That is still enormously better than the zero
numbers available before, and the per-route breakdown separates the worst
of it — but it should not be read as a physical constant.
The gate is also deliberately ASYMMETRIC. It only ever RAISES the bar, and
only once there are DECAY_MIN_SAMPLES real attempts to derive it from.
With no data it returns the static floor unchanged — this module cannot
make the bot trade more, only stop it attempting trades the data says are
already lost. `T_pipeline` uses the p90 of observed elapsed times rather
than the median, because the trade that matters is the slow one.
Set DECAY_GATE_ENABLED=false to record without gating — recommended for
the first day, so the numbers can be read (`/decay`) before they act.
"""
from __future__ import annotations
import csv
import logging
import os
import threading
import time
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Optional
logger = logging.getLogger(__name__)
def _env_float(name: str, default: float) -> float:
try:
raw = os.getenv(name, "").strip()
return float(raw) if raw else default
except (TypeError, ValueError):
return default
def _env_int(name: str, default: int) -> int:
try:
raw = os.getenv(name, "").strip()
return int(float(raw)) if raw else default
except (TypeError, ValueError):
return default
def _env_bool(name: str, default: bool) -> bool:
raw = os.getenv(name, "").strip().lower()
if not raw:
return default
return raw in ("1", "true", "yes", "on")
# Recording is always on — it is three floats per real attempt and it is the
# only way this deployment will ever stop guessing.
RECORD_ENABLED = _env_bool("DECAY_RECORD_ENABLED", True)
# Gating ships OFF. Read the numbers for a day first; a gate derived from a
# handful of samples is a guess with extra steps.
GATE_ENABLED = _env_bool("DECAY_GATE_ENABLED", False)
_MIN_SAMPLES = max(5, _env_int("DECAY_MIN_SAMPLES", 20))
_WINDOW = max(_MIN_SAMPLES, _env_int("DECAY_WINDOW_SAMPLES", 200))
# Age ceiling on a decay sample. A measurement of how much edge died in a
# 2-second pipeline says nothing about a 0.2-second one, and keeping it
# lets a rebuilt pipeline stay priced at its old cost forever — see
# _recent() for the loop that produced.
_MAX_AGE_SECS = max(3600.0, float(_env_int("DECAY_MAX_AGE_SECS", 86400)))
# How many of the newest samples count as "recent" for the trend split.
# 20 is small enough to move within an hour of real attempts and large
# enough that one outlier cannot define it.
_RECENT_N = max(5, _env_int("DECAY_RECENT_SAMPLES", 20))
# Cap the gate's contribution so a pathological sample can't lock trading
# out entirely — same reasoning as dynamic_floor's own multiplier cap.
_MAX_ADDER_USD = _env_float("DECAY_MAX_ADDER_USD", 5.0)
_PATH = os.getenv("DECAY_LOG_PATH", "").strip() or "data/decay_log.csv"
_COLUMNS = [
"ts_iso", "route", "quote_net_usd", "requote_net_usd", "elapsed_ms",
"decay_usd", "decay_usd_per_sec", "outcome", "detail",
]
@dataclass
class _Sample:
route: str
quote_net_usd: float
requote_net_usd: float
elapsed_ms: float
at: float = field(default_factory=time.time)
@property
def decay_usd(self) -> float:
return self.quote_net_usd - self.requote_net_usd
@property
def decay_per_sec(self) -> float:
secs = self.elapsed_ms / 1000.0
return self.decay_usd / secs if secs > 0.05 else 0.0
def _percentile(values: list[float], q: float) -> float:
if not values:
return 0.0
ordered = sorted(values)
idx = min(len(ordered) - 1, max(0, int(round(q * (len(ordered) - 1)))))
return ordered[idx]
def _median(values: list[float]) -> float:
return _percentile(values, 0.5)
class DecayTracker:
def __init__(self) -> None:
self._samples: list[_Sample] = []
self._lock = threading.Lock()
self._path = Path(_PATH).expanduser()
self._ready = False
if RECORD_ENABLED:
try:
self._path.parent.mkdir(parents=True, exist_ok=True)
if not self._path.exists() or self._path.stat().st_size == 0:
with self._path.open("w", newline="", encoding="utf-8") as fh:
csv.writer(fh).writerow(_COLUMNS)
self._ready = True
self._hydrate()
logger.info(
"[DecayTracker] recording attempts to %s (%d prior sample(s) "
"loaded)", self._path, len(self._samples),
)
except Exception as exc: # noqa: BLE001 — telemetry never blocks trading
logger.warning("[DecayTracker] could not open %s (%s)", self._path, exc)
def _hydrate(self) -> None:
"""Reload recent samples from the CSV at startup (v1.1).
Caught in review on PR #132: samples lived in memory only, so every
restart reset the estimator to zero. That is fatal for THIS
deployment specifically — the operator restarts on every config
change, often several times an hour, and real execution attempts
are rare. The window would have reset faster than it could ever
fill, so `/decay` would sit at "0/20 attempts" forever and the gate
could never engage no matter how long the bot ran.
Reads only the tail, tolerates a truncated first line, and skips
any row it cannot parse. A corrupt or missing file means "start
empty" — never a crash on startup.
"""
try:
size = self._path.stat().st_size
with self._path.open("rb") as fh:
# ~200 bytes/row; read generously and slice after parsing.
budget = _WINDOW * 400
if size > budget:
fh.seek(size - budget)
fh.readline() # drop the partial row the seek landed in
blob = fh.read().decode("utf-8", errors="replace")
except OSError as exc:
logger.debug("[DecayTracker] no prior samples to load: %s", exc)
return
loaded: list[_Sample] = []
for row in csv.reader(blob.splitlines()):
if len(row) != len(_COLUMNS) or row[0] == "ts_iso":
continue
try:
# v1.17 — CARRY THE ROW'S OWN TIMESTAMP.
#
# `at` defaults to time.time(), so loading rebuilt every
# historical sample as if it had just been measured. That
# alone would have defeated the age ceiling added in this
# version: the operator restarts several times an hour, and
# each restart would have stamped five-day-old measurements
# as fresh and handed them straight back to the gate.
#
# Column 0 is the ISO timestamp the row was written with.
# A row whose timestamp will not parse is dropped rather
# than dated to now — an undateable sample cannot be aged
# out, and one that cannot be aged out is exactly what this
# is here to prevent.
at = datetime.fromisoformat(
row[0].replace("Z", "+00:00")).timestamp()
loaded.append(_Sample(
route=row[1], quote_net_usd=float(row[2]),
requote_net_usd=float(row[3]), elapsed_ms=float(row[4]),
at=at,
))
except (TypeError, ValueError):
continue
if loaded:
self._samples = loaded[-_WINDOW:]
# ── Recording ────────────────────────────────────────────────────────
def record(
self, route: str, quote_net_usd: float, requote_net_usd: float,
elapsed_ms: float, outcome: str = "", detail: str = "",
) -> None:
"""One real execution attempt. Never raises."""
try:
sample = _Sample(
route=route, quote_net_usd=float(quote_net_usd),
requote_net_usd=float(requote_net_usd), elapsed_ms=float(elapsed_ms),
)
except (TypeError, ValueError):
return
with self._lock:
self._samples.append(sample)
if len(self._samples) > _WINDOW * 2:
del self._samples[: len(self._samples) - _WINDOW]
logger.info(
"[DecayTracker] %s — quote $%+.4f -> re-quote $%+.4f after %.0fms "
"(decay $%.4f, $%.4f/s) — %s",
route, sample.quote_net_usd, sample.requote_net_usd, sample.elapsed_ms,
sample.decay_usd, sample.decay_per_sec, outcome or "?",
)
if not self._ready:
return
try:
row = [
datetime.now(timezone.utc).isoformat(timespec="seconds"), route,
round(sample.quote_net_usd, 6), round(sample.requote_net_usd, 6),
round(sample.elapsed_ms, 1), round(sample.decay_usd, 6),
round(sample.decay_per_sec, 6), outcome, (detail or "")[:200],
]
with self._lock, self._path.open("a", newline="", encoding="utf-8") as fh:
csv.writer(fh).writerow(row)
except Exception as exc: # noqa: BLE001
logger.debug("[DecayTracker] row skipped: %s", exc)
# ── Estimation ───────────────────────────────────────────────────────
def _recent(self, route: Optional[str] = None) -> list[_Sample]:
"""The window, bounded by AGE as well as by count (v1.17).
────────────────────────────────────────────────────────────────
THE FOREVER-LOOP THIS CLOSES
────────────────────────────────────────────────────────────────
The window was 200 samples and nothing else. `_Sample.at` has
existed since v1.0 and was never read, so a measurement taken on a
pipeline that has since been rebuilt still set today's floor.
On the operator's box that became self-sustaining:
old 2.05s pipeline → decay measured at $0.92/s
→ gate charges $2.566 of headroom
→ a configured $0.05 floor becomes $2.62
→ nothing clears $2.62, so nothing sends
→ no new attempt is ever recorded
→ the 20 ancient samples remain the whole window, forever
Their /pipeline now reports 1.19s over 3 attempts while /decay
still reports 2.05s over 20 — the same shape as the flash-fee
loop fixed earlier: a stale safety number pricing itself into the
gate that prevents its own replacement.
Samples older than DECAY_MAX_AGE_SECS (24h) are dropped. If that
leaves too few, the caller sees an under-sampled window and the
gate stands down — which is correct: no recent evidence is a
reason to stop charging for decay, not a reason to keep charging
yesterday's rate.
"""
cutoff = time.time() - _MAX_AGE_SECS
with self._lock:
samples = [s for s in self._samples[-_WINDOW:]
if getattr(s, "at", 0.0) >= cutoff]
if route:
samples = [s for s in samples if s.route == route]
return samples
def decay_usd_per_sec(self, route: Optional[str] = None) -> Optional[float]:
"""Median decay rate, or None without a real sample.
Route-specific when there is enough of it, global otherwise: a
per-route estimate off three samples is noise wearing a label.
Negative medians are clamped to 0 — an edge that reliably GREW
between quote and re-quote is a sampling artefact, and betting on
it would be exactly the wrong direction to be wrong in.
"""
samples = self._recent(route)
if route and len(samples) < _MIN_SAMPLES:
samples = self._recent(None)
if len(samples) < _MIN_SAMPLES:
return None
rates = [s.decay_per_sec for s in samples if s.elapsed_ms > 50]
if not rates:
return None
return max(0.0, _median(rates))
def pipeline_secs_p90(self) -> Optional[float]:
samples = self._recent(None)
if len(samples) < _MIN_SAMPLES:
return None
return _percentile([s.elapsed_ms for s in samples], 0.90) / 1000.0
def _gate_samples(self, route: Optional[str] = None) -> list[_Sample]:
"""Samples the gate derives its adder from — RECENT first (v1.3).
The long window is right for describing history and wrong for
setting a bar, because a structural change to the pipeline takes
hours to work through 200 samples. Prefer the newest _RECENT_N when
that slice is itself large enough to be an estimate rather than an
anecdote; fall back to the full window otherwise.
"""
samples = self._recent(route)
if route and len(samples) < _MIN_SAMPLES:
samples = self._recent(None)
if len(samples) >= _RECENT_N >= _MIN_SAMPLES:
return samples[-_RECENT_N:]
return samples
def decay_usd_p90(self, route: Optional[str] = None) -> Optional[float]:
"""p90 of the dollars actually lost between quote and re-quote.
This is the number the gate wants, and it is measured directly.
Clamped at 0: a p90 decay below zero would mean the edge usually
GREW, which is not something to hand back as a discount.
"""
samples = self._gate_samples(route)
if len(samples) < _MIN_SAMPLES:
return None
return max(0.0, _percentile([s.decay_usd for s in samples], 0.90))
def required_floor(
self, static_floor_usd: float, route: Optional[str] = None,
) -> tuple[float, str]:
"""(floor_to_enforce, human explanation).
Implements the operator's catch condition:
Net_at_detection ≥ MIN_PROFIT_FLOOR + (edge lost during the pipeline)
v1.3 — the second term is now MEASURED, not reconstructed.
It used to be `decay_rate_$/s x pipeline_p90_secs`, and that product
broke the moment the pipeline got faster. Two independent reasons,
both visible in this deployment's own numbers on 2026-07-30:
THE RATE INFLATES AS THE PIPELINE SHRINKS. decay_per_sec is
decay_usd / elapsed_secs, and most of decay_usd is quote noise
that does NOT shrink with elapsed time — only the adversarial and
drift components do. So the same ~$1.6 of quote wobble reads as
$0.13/s over a 12s pipeline and $1.07/s over a 1.5s one. The
measured rate DOUBLED here at the exact moment the pipeline got
eight times faster. That is division, not the market.
THE TWO FACTORS CAME FROM DIFFERENT WORLDS. The rate is dominated
by new fast samples while the p90 elapsed still lags on 200 rows
of the old slow pipeline. Multiplying them produced $15.69 — an
adder describing a pipeline that no longer exists, saturating the
$5.00 cap and reporting the cap as if it were a finding. On edges
worth $2-5, an enabled gate would have refused literally every
trade, and the /decay output would have explained the refusal with
a number that was never observed.
The dataset already contains decay_usd per attempt. Reconstructing
it from a rate and a duration only makes sense when extrapolating to
a duration you have not observed — and here we have observed it. So
take the p90 of the dollars directly. It cannot decouple from the
pipeline, because it IS what the pipeline cost.
The rate is still computed and still reported: "how much does a
second cost me" is exactly the question that says whether more
latency work pays. It just no longer sets the bar.
"""
static = max(0.0, float(static_floor_usd))
if not GATE_ENABLED:
return static, "decay gate off (DECAY_GATE_ENABLED=false) — recording only"
observed = self.decay_usd_p90(route)
if observed is None:
return static, (
f"fewer than {_MIN_SAMPLES} real attempts recorded — no decay "
f"estimate yet, static floor unchanged"
)
adder = min(_MAX_ADDER_USD, observed)
n = len(self._gate_samples(route))
return static + adder, (
f"observed decay p90 ${observed:.4f} over the last {n} attempt(s) "
f"= ${adder:.4f} on top of ${static:.2f}"
+ (f" (CAPPED at ${_MAX_ADDER_USD:.2f})" if adder >= _MAX_ADDER_USD else "")
)
# ── Introspection (/decay) ───────────────────────────────────────────
def status(self) -> dict[str, Any]:
samples = self._recent(None)
by_route: dict[str, list[_Sample]] = {}
for s in samples:
by_route.setdefault(s.route, []).append(s)
rate = self.decay_usd_per_sec()
pipeline = self.pipeline_secs_p90()
# v1.2 — the RECENT slice, reported next to the full window.
#
# The window is 200 samples and hydrates from disk, so a change to
# the pipeline is invisible here for hours: a p90 over a window
# still dominated by pre-change samples reports the OLD pipeline
# and looks like the fix did nothing. It did that immediately after
# the execute-on-detection change — 14.64s became 15.63s, which
# reads as a regression and is really just the same old samples
# plus one.
#
# A metric that cannot show an improvement for hours is a metric
# nobody can act on. The last _RECENT_N samples are now summarised
# separately, so the effect of a change is visible within a handful
# of attempts and the long window keeps doing its job for the gate.
recent = samples[-_RECENT_N:] if samples else []
older = samples[:-_RECENT_N] if len(samples) > _RECENT_N else []
def _slice(rows: list[_Sample]) -> dict[str, Any]:
if not rows:
return {"samples": 0, "median_secs": None, "median_decay_usd": None}
return {
"samples": len(rows),
"median_secs": round(_median([s.elapsed_ms for s in rows]) / 1000.0, 3),
"median_decay_usd": round(_median([s.decay_usd for s in rows]), 4),
}
recent_stats, older_stats = _slice(recent), _slice(older)
trend = None
if recent_stats["median_secs"] is not None and older_stats["median_secs"]:
trend = round(
recent_stats["median_secs"] - older_stats["median_secs"], 3,
)
# v1.3 — is the "decay" directional, or is it noise?
#
# Decay is supposed to be adversarial: other searchers take the edge,
# so it should shrink far more often than it grows. If instead the
# edge grows on roughly half the attempts, then what is being
# measured is not decay at all — it is the ordinary variance between
# two Jupiter quotes taken moments apart, and the median of it is a
# number with no mechanism behind it.
#
# That distinction decides where the next hour of work goes. Skewed
# means latency work pays. Symmetric means latency work is finished
# and the constraint is that the edges are inside the noise band.
# Without this split both look identical in the median.
grew = sum(1 for s in samples if s.decay_usd < 0)
gate_rows = self._gate_samples(None)
observed_p90 = self.decay_usd_p90()
return {
"recording": self._ready,
"gate_enabled": GATE_ENABLED,
"path": str(self._path),
"samples": len(samples),
"min_samples": _MIN_SAMPLES,
# Newest first in the reading order that matters: "is it better
# NOW" before "what is the long-run average".
"recent": recent_stats,
"earlier": older_stats,
"pipeline_trend_secs": trend,
"recent_n": _RECENT_N,
"decay_usd_per_sec": None if rate is None else round(rate, 5),
# What the gate actually charges: dollars observed, not a rate
# multiplied by a duration from a different window.
"observed_decay_p90_usd": (
None if observed_p90 is None else round(observed_p90, 4)
),
"gate_window": len(gate_rows),
"gate_capped": (
observed_p90 is not None and observed_p90 >= _MAX_ADDER_USD
),
"max_adder_usd": _MAX_ADDER_USD,
# The noise check. `grew` counts attempts where the edge was
# BIGGER at send time than at detection. Adversarial decay
# cannot do that often; symmetric quote noise does it half the
# time. Near 50% means the median "decay" has no mechanism
# behind it and latency work will not move it.
"grew_count": grew,
"grew_share": round(grew / len(samples), 3) if samples else None,
# v1.4 — the plain median decay across all attempts, in dollars.
# Everything above reports a rate, a p90 or a per-route figure;
# none of them is the like-for-like number modules/decay_control
# needs to set the treatment arm against its control. Comparing
# a p90 to a median would make the control look flat no matter
# what it did.
"median_decay_usd": (
round(_median([s.decay_usd for s in samples]), 5) if samples else None
),
"pipeline_p90_secs": None if pipeline is None else round(pipeline, 3),
"pipeline_median_secs": (
round(_median([s.elapsed_ms for s in samples]) / 1000.0, 3)
if samples else None
),
# v1.3 — the adder the gate WOULD apply, from the same code path
# the gate uses. It used to be recomputed here as rate x p90,
# which meant /decay could report a number the gate would never
# charge. One source of truth.
"implied_adder_usd": (
None if observed_p90 is None
else round(min(_MAX_ADDER_USD, observed_p90), 4)
),
"routes": [
{
"route": route,
"attempts": len(rows),
"median_decay_usd": round(_median([s.decay_usd for s in rows]), 4),
"median_elapsed_ms": round(_median([s.elapsed_ms for s in rows])),
"survived": sum(1 for s in rows if s.requote_net_usd > 0),
}
for route, rows in sorted(
by_route.items(), key=lambda kv: -len(kv[1]),
)[:8]
],
}
_tracker: Optional[DecayTracker] = None
def get_tracker() -> DecayTracker:
global _tracker
if _tracker is None:
_tracker = DecayTracker()
return _tracker