diff --git a/prisma.config.ts b/prisma.config.ts new file mode 100644 index 00000000..a46709cf --- /dev/null +++ b/prisma.config.ts @@ -0,0 +1,11 @@ +import { defineConfig, env } from "prisma/config"; + +export default defineConfig({ + schema: "prisma/schema.prisma", + migrations: { + path: "prisma/migrations", + }, + datasource: { + url: env("DATABASE_URL"), + }, +}); diff --git a/src/network/nonce_tracker.py b/src/network/nonce_tracker.py index 361bbaca..54911fb2 100644 --- a/src/network/nonce_tracker.py +++ b/src/network/nonce_tracker.py @@ -1,9 +1,10 @@ import asyncio import logging +import math import os import threading import time -from collections import defaultdict +from collections import defaultdict, deque from dataclasses import dataclass, field from typing import Any, Callable, Dict, List, Optional, Set @@ -21,6 +22,126 @@ # get_stale(address, timeout_seconds=...). DEFAULT_STALE_TIMEOUT_SECONDS = 30.0 +# Nominal Stellar ledger close interval (~5 s). Used to translate observed +# consensus closing latency into a ledger-index buffer for Soroban submissions. +DEFAULT_LEDGER_CLOSE_SECONDS = 5.0 +DEFAULT_CONSENSUS_CLOSING_LATENCY_SECONDS = 5.0 +DEFAULT_MIN_LEDGER_BUFFER = 2 +DEFAULT_MAX_LEDGER_BUFFER = 64 +DEFAULT_LEDGER_SAFETY_MARGIN = 1 +CONSENSUS_CLOSING_LATENCY_WINDOW = MOVING_AVG_WINDOW_SIZE + + +@dataclass(frozen=True) +class LedgerTimeBounds: + """Inclusive ledger sequence window for Stellar transaction timebounds.""" + + min_ledger: int + max_ledger: int + + +class AdaptiveLedgerTimeBoundCalculator: + """Adaptive min/max ledger bounds for Soroban submission envelopes. + + Observed consensus *closing* latency (how long recent ledgers took to + close) is tracked in a bounded moving average. The ledger buffer applied + to ``max_ledger`` scales with that average so submissions remain valid + when the network is slow to reach consensus. + + Parameters + ---------- + ledger_close_seconds: + Expected ledger close duration under nominal conditions. + min_ledger_buffer: + Minimum number of ledgers between ``min_ledger`` and ``max_ledger``. + max_ledger_buffer: + Hard cap on the adaptive ledger span. + safety_margin_ledgers: + Extra ledgers added on top of the latency-derived estimate. + latency_window: + Number of recent closing-latency samples in the moving average. + default_closing_latency_seconds: + Assumed closing latency before any samples are recorded. + """ + + def __init__( + self, + ledger_close_seconds: float = DEFAULT_LEDGER_CLOSE_SECONDS, + min_ledger_buffer: int = DEFAULT_MIN_LEDGER_BUFFER, + max_ledger_buffer: int = DEFAULT_MAX_LEDGER_BUFFER, + safety_margin_ledgers: int = DEFAULT_LEDGER_SAFETY_MARGIN, + latency_window: int = CONSENSUS_CLOSING_LATENCY_WINDOW, + default_closing_latency_seconds: float = DEFAULT_CONSENSUS_CLOSING_LATENCY_SECONDS, + ) -> None: + if ledger_close_seconds <= 0: + raise ValueError("ledger_close_seconds must be positive") + if min_ledger_buffer < 1: + raise ValueError("min_ledger_buffer must be >= 1") + if max_ledger_buffer < min_ledger_buffer: + raise ValueError("max_ledger_buffer must be >= min_ledger_buffer") + if safety_margin_ledgers < 0: + raise ValueError("safety_margin_ledgers must be >= 0") + if latency_window < 1: + raise ValueError("latency_window must be >= 1") + if default_closing_latency_seconds <= 0: + raise ValueError("default_closing_latency_seconds must be positive") + + self._ledger_close_seconds = ledger_close_seconds + self._min_ledger_buffer = min_ledger_buffer + self._max_ledger_buffer = max_ledger_buffer + self._safety_margin_ledgers = safety_margin_ledgers + self._latency_window = latency_window + self._default_closing_latency_seconds = default_closing_latency_seconds + self._closing_latency_samples: deque[float] = deque(maxlen=latency_window) + self._lock = threading.Lock() + + @property + def min_ledger_buffer(self) -> int: + return self._min_ledger_buffer + + def record_consensus_closing_latency(self, latency_seconds: float) -> None: + """Record how long a ledger took to close (consensus round-trip).""" + if latency_seconds <= 0: + raise ValueError("latency_seconds must be positive") + with self._lock: + self._closing_latency_samples.append(float(latency_seconds)) + + @property + def consensus_closing_latency_seconds(self) -> float: + """Moving average of recorded closing latencies, or the default.""" + with self._lock: + if not self._closing_latency_samples: + return self._default_closing_latency_seconds + return sum(self._closing_latency_samples) / len( + self._closing_latency_samples + ) + + def _adaptive_ledger_buffer(self) -> int: + latency = self.consensus_closing_latency_seconds + estimated = math.ceil(latency / self._ledger_close_seconds) + estimated += self._safety_margin_ledgers + return max( + self._min_ledger_buffer, + min(estimated, self._max_ledger_buffer), + ) + + def compute_bounds(self, current_ledger: int) -> LedgerTimeBounds: + """Return ledger timebounds anchored at *current_ledger*. + + ``min_ledger`` is the current ledger sequence. ``max_ledger`` extends + forward by a buffer derived from consensus closing latency. + """ + if current_ledger < 0: + raise ValueError("current_ledger must be >= 0") + buffer = self._adaptive_ledger_buffer() + return LedgerTimeBounds( + min_ledger=current_ledger, + max_ledger=current_ledger + buffer, + ) + + +adaptive_ledger_time_bound_calculator = AdaptiveLedgerTimeBoundCalculator() + class HorizonNodeProfile: def __init__(self, name: str, url: str): @@ -1302,6 +1423,8 @@ def _run_monitor(self) -> None: __all__ = [ + "AdaptiveLedgerTimeBoundCalculator", + "LedgerTimeBounds", "NonceTracker", "NonceWindow", "NonceGapDetector", @@ -1312,6 +1435,7 @@ def _run_monitor(self) -> None: "ReconciliationResult", "nonce_tracker", "nonce_window", + "adaptive_ledger_time_bound_calculator", "RPCNodeFailoverSupervisor", "rpc_supervisor", "PredictiveRPCSupervisor", diff --git a/tests/test_nonce_tracker.py b/tests/test_nonce_tracker.py index b007d1a4..a3e32d98 100644 --- a/tests/test_nonce_tracker.py +++ b/tests/test_nonce_tracker.py @@ -12,7 +12,9 @@ sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "src")) from network.nonce_tracker import ( + AdaptiveLedgerTimeBoundCalculator, GapReport, + LedgerTimeBounds, NonceGapDetector, NonceRecoveryEngine, NonceTracker, @@ -1141,3 +1143,40 @@ def test_transaction_latency_tracking(caplog: pytest.LogCaptureFixture) -> None: assert "latency=" in caplog.text assert "latency=0.0ms" not in caplog.text + +# =========================================================================== +# Adaptive ledger time-bounds — Issue #651 +# =========================================================================== + + +def test_adaptive_time_bounds() -> None: + """Acceptance test for Issue #651: bounds account for consensus closing latency.""" + calc = AdaptiveLedgerTimeBoundCalculator( + ledger_close_seconds=5.0, + min_ledger_buffer=2, + safety_margin_ledgers=1, + ) + current_ledger = 123_456 + + baseline = calc.compute_bounds(current_ledger=current_ledger) + assert isinstance(baseline, LedgerTimeBounds) + assert baseline.min_ledger == current_ledger + assert baseline.max_ledger == current_ledger + calc.min_ledger_buffer + + for _ in range(4): + calc.record_consensus_closing_latency(20.0) + + under_load = calc.compute_bounds(current_ledger=current_ledger) + assert under_load.max_ledger > baseline.max_ledger + + fast = AdaptiveLedgerTimeBoundCalculator( + ledger_close_seconds=5.0, + min_ledger_buffer=2, + safety_margin_ledgers=1, + ) + for _ in range(4): + fast.record_consensus_closing_latency(2.5) + + nominal = fast.compute_bounds(current_ledger=current_ledger) + assert nominal.max_ledger <= under_load.max_ledger + assert nominal.max_ledger >= current_ledger + fast.min_ledger_buffer