Skip to content

Commit 7e39531

Browse files
HenryHenry
authored andcommitted
fix(runtime): harden incremental candle caching
1 parent 74ce431 commit 7e39531

5 files changed

Lines changed: 359 additions & 36 deletions

File tree

backend_api_python/app/data_sources/crypto.py

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -698,6 +698,14 @@ def get_kline(
698698
or int(klines[-1]["time"]) < int(before_time) - tolerance_seconds
699699
or coverage_ratio < 0.98
700700
):
701+
self._set_last_failure(
702+
"Incomplete K-line coverage after normalization: "
703+
f"requested={after_time}~{before_time}, "
704+
f"actual={klines[0]['time']}~{klines[-1]['time']}, "
705+
f"rows={len(klines)}/{expected_rows}",
706+
symbol=symbol_pair,
707+
timeframe=timeframe,
708+
)
701709
logger.warning(
702710
"Rejected incomplete %s %s K-lines after normalization: "
703711
"requested=%s~%s, actual=%s~%s, rows=%s/%s (%.2f%%)",
@@ -982,6 +990,13 @@ def _fetch_ohlcv(
982990
int(ohlcv[0][0]) > requested_start_ms + tolerance_ms
983991
or int(ohlcv[-1][0]) < end_ms - tolerance_ms
984992
):
993+
self._set_last_failure(
994+
"Incomplete K-line history: "
995+
f"requested={requested_start_ms}~{end_ms}, "
996+
f"actual={int(ohlcv[0][0])}~{int(ohlcv[-1][0])}",
997+
symbol=symbol_pair,
998+
timeframe=timeframe,
999+
)
9851000
logger.warning(
9861001
"Refused incomplete %s %s history: requested=%s~%s, actual=%s~%s",
9871002
exchange_id,

backend_api_python/app/data_sources/errors.py

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,19 @@ def classify_market_data_failure(
9595
code = "rate_limited"
9696
message = "The exchange rate limit was reached. Market data will be retried."
9797
retryable = True
98+
elif any(token in text for token in (
99+
"incomplete k-line",
100+
"incomplete kline",
101+
"incomplete candle",
102+
"incomplete market data",
103+
"incomplete history",
104+
)):
105+
code = "incomplete_market_data"
106+
message = (
107+
"The exchange returned incomplete K-line coverage. "
108+
"The missing interval will be retried."
109+
)
110+
retryable = True
98111
elif any(token in text for token in ("timeout", "timed out", "network error", "connection reset", "connection refused", "exchange not available", "service unavailable", "502", "503", "504")):
99112
code = "exchange_unavailable"
100113
message = "The exchange market-data service is temporarily unreachable."

backend_api_python/app/services/strategy_v2/market_data.py

Lines changed: 153 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,10 @@
1111
import pandas as pd
1212

1313
from app.data_sources import DataSourceFactory
14-
from app.data_sources.errors import MarketDataUnavailableError
14+
from app.data_sources.errors import (
15+
MarketDataUnavailableError,
16+
classify_market_data_failure,
17+
)
1518
from app.services.backtest_cache import KlineCache
1619
from app.utils.logger import get_logger
1720

@@ -30,6 +33,8 @@ class _SharedFrameEntry:
3033
_shared_frame_locks: dict[str, threading.RLock] = {}
3134
_shared_frames_lock = threading.RLock()
3235
_SHARED_FRAME_CACHE_MAX_SIZE = 256
36+
_LIVE_CACHE_GRACE_BARS = 2
37+
_INCOMPLETE_WARMUP_RETRIES = 1
3338

3439
TIMEFRAME_SECONDS = {
3540
"1m": 60,
@@ -58,6 +63,25 @@ def _normalize_utc_datetime(value: datetime) -> datetime:
5863
return value.astimezone(timezone.utc)
5964

6065

66+
def _last_completed_bar_open(
67+
timeframe_seconds: int,
68+
*,
69+
now: Optional[datetime] = None,
70+
) -> pd.Timestamp:
71+
"""Return the UTC open time of the most recently completed candle.
72+
73+
Using a bar-aligned cutoff avoids treating the seconds inside the current
74+
minute as an uncovered cache tail and refetching the same 1m candle on
75+
every runtime heartbeat.
76+
"""
77+
seconds = max(1, int(timeframe_seconds or 1))
78+
current = _normalize_utc_datetime(now or datetime.now(timezone.utc))
79+
completed_open = ((int(current.timestamp()) // seconds) - 1) * seconds
80+
return pd.Timestamp(
81+
datetime.fromtimestamp(completed_open, tz=timezone.utc)
82+
).tz_localize(None)
83+
84+
6185
def _covers_crypto_window(
6286
frame: pd.DataFrame,
6387
requested_start: pd.Timestamp,
@@ -110,8 +134,8 @@ def _load_strategy_frame_uncached(
110134
))
111135
requested_start = pd.Timestamp(start_utc).tz_localize(None)
112136
requested_end = pd.Timestamp(end_utc).tz_localize(None)
113-
closed_bar_cutoff = datetime.now(timezone.utc).replace(tzinfo=None) - timedelta(seconds=timeframe_seconds)
114-
coverage_end = min(requested_end, pd.Timestamp(closed_bar_cutoff))
137+
closed_bar_cutoff = _last_completed_bar_open(timeframe_seconds)
138+
coverage_end = min(requested_end, closed_bar_cutoff)
115139
cached = _cache.get(cache_key)
116140
if cached is not None and not cached.empty:
117141
if str(market or "").strip().lower() != "crypto" or _covers_crypto_window(
@@ -187,7 +211,7 @@ def _load_strategy_frame_uncached(
187211
subset=["open", "high", "low", "close"]
188212
)
189213
if requested_end >= closed_bar_cutoff:
190-
frame = frame[frame.index <= pd.Timestamp(closed_bar_cutoff)]
214+
frame = frame[frame.index <= closed_bar_cutoff]
191215
if (
192216
str(market or "").strip().lower() == "crypto"
193217
and not _covers_crypto_window(frame, requested_start, coverage_end, timeframe_seconds)
@@ -202,7 +226,15 @@ def _load_strategy_frame_uncached(
202226
frame.index.min(),
203227
frame.index.max(),
204228
)
205-
return pd.DataFrame()
229+
raise MarketDataUnavailableError(
230+
classify_market_data_failure(
231+
"Incomplete K-line coverage after strategy window normalization",
232+
exchange_id=exchange_id or "",
233+
market_type=market_type or "",
234+
symbol=symbol,
235+
timeframe=timeframe,
236+
)
237+
)
206238
if not frame.empty:
207239
_cache.put(cache_key, frame, timeframe)
208240
return frame.copy()
@@ -253,6 +285,38 @@ def _evict_shared_frame_if_needed() -> None:
253285
_shared_frames.pop(oldest_key, None)
254286

255287

288+
def _is_live_request(
289+
requested_end: pd.Timestamp,
290+
closed_cutoff: pd.Timestamp,
291+
timeframe_seconds: int,
292+
) -> bool:
293+
tolerance = pd.Timedelta(seconds=max(1, timeframe_seconds) * 2)
294+
return bool(requested_end >= closed_cutoff - tolerance)
295+
296+
297+
def _cached_crypto_frame_is_usable(
298+
frame: pd.DataFrame,
299+
requested_start: pd.Timestamp,
300+
coverage_end: pd.Timestamp,
301+
timeframe_seconds: int,
302+
*,
303+
live_request: bool,
304+
) -> bool:
305+
if not _covers_crypto_window(
306+
frame,
307+
requested_start,
308+
coverage_end,
309+
timeframe_seconds,
310+
):
311+
return False
312+
if not live_request:
313+
return True
314+
max_lag = pd.Timedelta(
315+
seconds=max(1, timeframe_seconds) * _LIVE_CACHE_GRACE_BARS
316+
)
317+
return bool(frame.index.max() >= coverage_end - max_lag)
318+
319+
256320
def clear_shared_strategy_frame_cache() -> None:
257321
"""Clear process-local candle state. Intended for tests and controlled reloads."""
258322
with _shared_frames_lock:
@@ -283,11 +347,14 @@ def load_strategy_frame(
283347
timeframe_seconds = TIMEFRAME_SECONDS.get(normalized_timeframe, 86400)
284348
requested_start = pd.Timestamp(start_utc).tz_localize(None)
285349
requested_end = pd.Timestamp(end_utc).tz_localize(None)
286-
closed_cutoff = pd.Timestamp(
287-
datetime.now(timezone.utc).replace(tzinfo=None)
288-
- timedelta(seconds=timeframe_seconds)
289-
)
350+
closed_cutoff = _last_completed_bar_open(timeframe_seconds)
290351
coverage_end = min(requested_end, closed_cutoff)
352+
live_request = _is_live_request(
353+
requested_end,
354+
closed_cutoff,
355+
timeframe_seconds,
356+
)
357+
crypto_market = str(market or "").strip().lower() == "crypto"
291358
key = _shared_frame_key(market, symbol, timeframe, market_type, exchange_id)
292359

293360
with _lock_for_shared_frame(key):
@@ -316,39 +383,98 @@ def load_strategy_frame(
316383
))
317384

318385
merged = entry.frame.copy() if entry is not None else pd.DataFrame()
319-
successful_windows: list[tuple[pd.Timestamp, pd.Timestamp]] = []
386+
last_failure: Optional[MarketDataUnavailableError] = None
320387
for window_start, window_end in fetch_windows:
321-
incoming = _load_strategy_frame_uncached(
322-
market,
323-
symbol,
324-
timeframe,
325-
window_start.to_pydatetime().replace(tzinfo=timezone.utc),
326-
window_end.to_pydatetime().replace(tzinfo=timezone.utc),
327-
market_type=market_type,
328-
exchange_id=exchange_id,
388+
incoming = pd.DataFrame()
389+
attempts = (
390+
1 + _INCOMPLETE_WARMUP_RETRIES
391+
if normalized_timeframe == "1m" and entry is None
392+
else 1
329393
)
394+
for attempt in range(attempts):
395+
try:
396+
incoming = _load_strategy_frame_uncached(
397+
market,
398+
symbol,
399+
timeframe,
400+
window_start.to_pydatetime().replace(tzinfo=timezone.utc),
401+
window_end.to_pydatetime().replace(tzinfo=timezone.utc),
402+
market_type=market_type,
403+
exchange_id=exchange_id,
404+
)
405+
last_failure = None
406+
break
407+
except MarketDataUnavailableError as exc:
408+
last_failure = exc
409+
should_retry = bool(
410+
exc.failure.code == "incomplete_market_data"
411+
and attempt + 1 < attempts
412+
)
413+
if should_retry:
414+
logger.warning(
415+
"Retrying incomplete %s %s warmup (%s/%s)",
416+
symbol,
417+
timeframe,
418+
attempt + 2,
419+
attempts,
420+
)
421+
continue
422+
break
330423
if incoming is not None and not incoming.empty:
331424
merged = _merge_frames(merged, incoming)
332-
successful_windows.append((window_start, min(window_end, coverage_end)))
333425

334426
if merged.empty:
427+
if last_failure is not None:
428+
raise last_failure
335429
return pd.DataFrame()
336430

337-
new_start = entry.coverage_start if entry is not None else requested_start
338-
new_end = entry.coverage_end if entry is not None else coverage_end
339-
for window_start, window_end in successful_windows:
340-
new_start = min(new_start, window_start)
341-
new_end = max(new_end, window_end)
431+
# Keep the cache bounded around the active warmup window. If a future
432+
# caller requests older history, the missing prefix is fetched again.
433+
overlap = pd.Timedelta(seconds=timeframe_seconds * 2)
434+
merged = merged[merged.index >= requested_start - overlap]
435+
actual_start = merged.index.min()
436+
actual_end = merged.index.max()
342437
_shared_frames[key] = _SharedFrameEntry(
343438
frame=merged,
344-
coverage_start=new_start,
345-
coverage_end=new_end,
439+
coverage_start=actual_start,
440+
coverage_end=actual_end,
346441
)
347442
_evict_shared_frame_if_needed()
348-
return merged[
443+
result = merged[
349444
(merged.index >= requested_start)
350445
& (merged.index <= coverage_end)
351446
].copy()
447+
if crypto_market and not _cached_crypto_frame_is_usable(
448+
result,
449+
requested_start,
450+
coverage_end,
451+
timeframe_seconds,
452+
live_request=live_request,
453+
):
454+
logger.warning(
455+
"Refused stale/incomplete cached crypto frame for %s %s: "
456+
"requested=%s~%s, actual=%s~%s",
457+
symbol,
458+
timeframe,
459+
requested_start,
460+
coverage_end,
461+
result.index.min() if not result.empty else "empty",
462+
result.index.max() if not result.empty else "empty",
463+
)
464+
if last_failure is not None:
465+
raise last_failure
466+
return pd.DataFrame()
467+
if last_failure is not None:
468+
logger.warning(
469+
"Using recent cached %s %s candles after transient %s failure; "
470+
"latest=%s, required=%s",
471+
symbol,
472+
timeframe,
473+
last_failure.failure.code,
474+
result.index.max(),
475+
coverage_end,
476+
)
477+
return result
352478

353479

354480
__all__ = [

backend_api_python/tests/test_market_data_failure_contract.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
("ProxyError: tunnel connection failed", "proxy_failure"),
2222
("binance does not have market symbol AVAX/USDT:USDT", "symbol_not_found"),
2323
("HTTP 429 too many requests", "rate_limited"),
24+
("Incomplete K-line coverage", "incomplete_market_data"),
2425
("request timed out", "exchange_unavailable"),
2526
],
2627
)

0 commit comments

Comments
 (0)