Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions prisma.config.ts
Original file line number Diff line number Diff line change
@@ -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"),
},
});
126 changes: 125 additions & 1 deletion src/network/nonce_tracker.py
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -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):
Expand Down Expand Up @@ -1302,6 +1423,8 @@ def _run_monitor(self) -> None:


__all__ = [
"AdaptiveLedgerTimeBoundCalculator",
"LedgerTimeBounds",
"NonceTracker",
"NonceWindow",
"NonceGapDetector",
Expand All @@ -1312,6 +1435,7 @@ def _run_monitor(self) -> None:
"ReconciliationResult",
"nonce_tracker",
"nonce_window",
"adaptive_ledger_time_bound_calculator",
"RPCNodeFailoverSupervisor",
"rpc_supervisor",
"PredictiveRPCSupervisor",
Expand Down
39 changes: 39 additions & 0 deletions tests/test_nonce_tracker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Loading