Fix three correctness bugs in the Yahoo ingest
Browse filesAll three are silent: nothing raises, the numbers are just wrong.
Unadjusted prices. yahoo.auto_adjust defaulted to false and _normalize_ohlcv
selects only the required columns, discarding 'adj close'. So splits were not
adjusted for: NVDA's June 2024 10:1 reads as a -90% single-day return and AAPL's
2020 4:1 as -75%. Any strategy or risk control spanning a split saw a phantom
crash. Default is now true, opting out logs a warning, and _warn_if_unadjusted
names the dates of daily-or-slower bars moving more than 35% as a backstop when
adjustment is off.
Partial bars reported as final. _poll_once pulls period='5d' and
_ingest_new_bars took every row past the watermark, including the period Yahoo
is still building. It emitted that in-progress bar, advanced _last_bar_ts past
it, and so never re-emitted the finished version — the strategy acted on a close
that had not happened yet, then never saw the real one. _drop_incomplete now
withholds a bar until its interval has elapsed; yahoo.emit_incomplete_bars
re-enables the old behaviour explicitly.
No backoff on rate limiting. The poll loop caught exceptions and retried at a
fixed 60s, so once Yahoo throttled us we stayed throttled. _poll_once now
reports success, and _next_delay backs off exponentially to
yahoo.max_backoff_seconds with jitter so symbols do not resynchronise after an
outage. connect() seeds the counter from its first attempt.
17 new tests, all of which fail against the previous implementation, including a
split fixture modelled on NVDA. Existing tests unchanged and still pass.
- agentic_ai_system/yahoo_data_stream.py +122 -12
- config.yaml +7 -1
- tests/test_yahoo_data_stream.py +169 -0
|
@@ -1,4 +1,5 @@
|
|
| 1 |
import logging
|
|
|
|
| 2 |
import threading
|
| 3 |
import time
|
| 4 |
from typing import Any, Callable, Dict, List, Optional
|
|
@@ -41,6 +42,26 @@ _MAX_LOOKBACK = {
|
|
| 41 |
'3mo': None,
|
| 42 |
}
|
| 43 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 44 |
|
| 45 |
class YahooDataStream:
|
| 46 |
"""
|
|
@@ -62,8 +83,13 @@ class YahooDataStream:
|
|
| 62 |
self.symbols = ['AAPL']
|
| 63 |
yahoo_cfg = config.get('yahoo', {})
|
| 64 |
self.poll_interval = int(yahoo_cfg.get('poll_interval_seconds', 60))
|
| 65 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 66 |
self.interval = self._map_interval(config.get('trading', {}).get('timeframe', '1d'))
|
|
|
|
| 67 |
self.data_callbacks: List[Callable] = []
|
| 68 |
self.is_connected = False
|
| 69 |
self.data_buffer: Dict[str, Dict[str, Any]] = {}
|
|
@@ -80,11 +106,21 @@ class YahooDataStream:
|
|
| 80 |
'latest_bar': None,
|
| 81 |
}
|
| 82 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 83 |
logger.info(
|
| 84 |
-
"Initialized YahooDataStream symbols=%s interval=%s poll_interval=%ss"
|
|
|
|
| 85 |
self.symbols,
|
| 86 |
self.interval,
|
| 87 |
self.poll_interval,
|
|
|
|
|
|
|
| 88 |
)
|
| 89 |
|
| 90 |
@staticmethod
|
|
@@ -102,7 +138,9 @@ class YahooDataStream:
|
|
| 102 |
return
|
| 103 |
|
| 104 |
self._stop_event.clear()
|
| 105 |
-
|
|
|
|
|
|
|
| 106 |
self._poll_thread = threading.Thread(target=self._poll_loop, name='yahoo-poll', daemon=True)
|
| 107 |
self._poll_thread.start()
|
| 108 |
self.is_connected = True
|
|
@@ -137,7 +175,8 @@ class YahooDataStream:
|
|
| 137 |
start, end = self._clamp_window(start_date, end_date, self.interval)
|
| 138 |
try:
|
| 139 |
raw = self._download(symbol, start=start, end=end, interval=self.interval)
|
| 140 |
-
df = self._normalize_ohlcv(raw)
|
|
|
|
| 141 |
if df.empty:
|
| 142 |
logger.warning("No Yahoo historical data for %s between %s and %s", symbol, start, end)
|
| 143 |
else:
|
|
@@ -173,8 +212,6 @@ class YahooDataStream:
|
|
| 173 |
}
|
| 174 |
|
| 175 |
def generate_simulated_data(self, symbol: str) -> Dict[str, Any]:
|
| 176 |
-
import random
|
| 177 |
-
|
| 178 |
latest_data = self.get_latest_data(symbol)
|
| 179 |
base_price = 150.0
|
| 180 |
if latest_data.get('latest_bar'):
|
|
@@ -197,13 +234,40 @@ class YahooDataStream:
|
|
| 197 |
return simulated_bar
|
| 198 |
|
| 199 |
def _poll_loop(self) -> None:
|
| 200 |
-
|
|
|
|
| 201 |
try:
|
| 202 |
-
self._poll_once()
|
| 203 |
except Exception as e:
|
| 204 |
logger.error("Yahoo poll loop error: %s", e, exc_info=True)
|
| 205 |
-
|
| 206 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 207 |
for symbol in self.symbols:
|
| 208 |
try:
|
| 209 |
raw = self._download(symbol, period='5d', interval=self.interval)
|
|
@@ -212,14 +276,60 @@ class YahooDataStream:
|
|
| 212 |
logger.warning("Yahoo poll returned no bars for %s", symbol)
|
| 213 |
continue
|
| 214 |
self._ingest_new_bars(symbol, df)
|
|
|
|
| 215 |
except Exception as e:
|
| 216 |
logger.error("Yahoo poll failed for %s: %s", symbol, e)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 217 |
|
| 218 |
def _ingest_new_bars(self, symbol: str, df: pd.DataFrame) -> None:
|
|
|
|
| 219 |
last_ts = self._last_bar_ts.get(symbol)
|
| 220 |
-
rows = df
|
| 221 |
if last_ts is not None:
|
| 222 |
-
rows =
|
| 223 |
if rows.empty:
|
| 224 |
return
|
| 225 |
|
|
|
|
| 1 |
import logging
|
| 2 |
+
import random
|
| 3 |
import threading
|
| 4 |
import time
|
| 5 |
from typing import Any, Callable, Dict, List, Optional
|
|
|
|
| 42 |
'3mo': None,
|
| 43 |
}
|
| 44 |
|
| 45 |
+
# How long each bar covers. Used to tell a finished bar from the one still
|
| 46 |
+
# forming right now -- see _drop_incomplete.
|
| 47 |
+
_INTERVAL_DURATION = {
|
| 48 |
+
'1m': pd.Timedelta(minutes=1),
|
| 49 |
+
'2m': pd.Timedelta(minutes=2),
|
| 50 |
+
'5m': pd.Timedelta(minutes=5),
|
| 51 |
+
'15m': pd.Timedelta(minutes=15),
|
| 52 |
+
'30m': pd.Timedelta(minutes=30),
|
| 53 |
+
'60m': pd.Timedelta(hours=1),
|
| 54 |
+
'90m': pd.Timedelta(minutes=90),
|
| 55 |
+
'1h': pd.Timedelta(hours=1),
|
| 56 |
+
'1d': pd.Timedelta(days=1),
|
| 57 |
+
'5d': pd.Timedelta(days=5),
|
| 58 |
+
'1wk': pd.Timedelta(weeks=1),
|
| 59 |
+
}
|
| 60 |
+
|
| 61 |
+
# A daily-or-slower bar that moves more than this is almost always an
|
| 62 |
+
# unadjusted split rather than a real move (NVDA's 2024 10:1 shows up as -90%).
|
| 63 |
+
_SPLIT_SUSPECT_MOVE = 0.35
|
| 64 |
+
|
| 65 |
|
| 66 |
class YahooDataStream:
|
| 67 |
"""
|
|
|
|
| 83 |
self.symbols = ['AAPL']
|
| 84 |
yahoo_cfg = config.get('yahoo', {})
|
| 85 |
self.poll_interval = int(yahoo_cfg.get('poll_interval_seconds', 60))
|
| 86 |
+
# Adjusted by default. With auto_adjust off, Yahoo returns raw Close and
|
| 87 |
+
# every split reads as a crash: NVDA's June 2024 10:1 becomes a -90% bar.
|
| 88 |
+
self.auto_adjust = bool(yahoo_cfg.get('auto_adjust', True))
|
| 89 |
+
self.emit_incomplete_bars = bool(yahoo_cfg.get('emit_incomplete_bars', False))
|
| 90 |
+
self.max_backoff = int(yahoo_cfg.get('max_backoff_seconds', 900))
|
| 91 |
self.interval = self._map_interval(config.get('trading', {}).get('timeframe', '1d'))
|
| 92 |
+
self._consecutive_failures = 0
|
| 93 |
self.data_callbacks: List[Callable] = []
|
| 94 |
self.is_connected = False
|
| 95 |
self.data_buffer: Dict[str, Dict[str, Any]] = {}
|
|
|
|
| 106 |
'latest_bar': None,
|
| 107 |
}
|
| 108 |
|
| 109 |
+
if not self.auto_adjust:
|
| 110 |
+
logger.warning(
|
| 111 |
+
"yahoo.auto_adjust is false: prices are NOT split- or dividend-adjusted. "
|
| 112 |
+
"Every split will appear as a large single-bar loss and any backtest "
|
| 113 |
+
"spanning one will be wrong."
|
| 114 |
+
)
|
| 115 |
+
|
| 116 |
logger.info(
|
| 117 |
+
"Initialized YahooDataStream symbols=%s interval=%s poll_interval=%ss "
|
| 118 |
+
"auto_adjust=%s emit_incomplete_bars=%s",
|
| 119 |
self.symbols,
|
| 120 |
self.interval,
|
| 121 |
self.poll_interval,
|
| 122 |
+
self.auto_adjust,
|
| 123 |
+
self.emit_incomplete_bars,
|
| 124 |
)
|
| 125 |
|
| 126 |
@staticmethod
|
|
|
|
| 138 |
return
|
| 139 |
|
| 140 |
self._stop_event.clear()
|
| 141 |
+
# Seed the backoff from the first attempt: if we are already being
|
| 142 |
+
# throttled, the loop should start backed off rather than hammering.
|
| 143 |
+
self._consecutive_failures = 0 if self._poll_once() else 1
|
| 144 |
self._poll_thread = threading.Thread(target=self._poll_loop, name='yahoo-poll', daemon=True)
|
| 145 |
self._poll_thread.start()
|
| 146 |
self.is_connected = True
|
|
|
|
| 175 |
start, end = self._clamp_window(start_date, end_date, self.interval)
|
| 176 |
try:
|
| 177 |
raw = self._download(symbol, start=start, end=end, interval=self.interval)
|
| 178 |
+
df = self._drop_incomplete(self._normalize_ohlcv(raw))
|
| 179 |
+
self._warn_if_unadjusted(symbol, df)
|
| 180 |
if df.empty:
|
| 181 |
logger.warning("No Yahoo historical data for %s between %s and %s", symbol, start, end)
|
| 182 |
else:
|
|
|
|
| 212 |
}
|
| 213 |
|
| 214 |
def generate_simulated_data(self, symbol: str) -> Dict[str, Any]:
|
|
|
|
|
|
|
| 215 |
latest_data = self.get_latest_data(symbol)
|
| 216 |
base_price = 150.0
|
| 217 |
if latest_data.get('latest_bar'):
|
|
|
|
| 234 |
return simulated_bar
|
| 235 |
|
| 236 |
def _poll_loop(self) -> None:
|
| 237 |
+
delay = self.poll_interval
|
| 238 |
+
while not self._stop_event.wait(delay):
|
| 239 |
try:
|
| 240 |
+
succeeded = self._poll_once()
|
| 241 |
except Exception as e:
|
| 242 |
logger.error("Yahoo poll loop error: %s", e, exc_info=True)
|
| 243 |
+
succeeded = False
|
| 244 |
+
self._consecutive_failures = 0 if succeeded else self._consecutive_failures + 1
|
| 245 |
+
delay = self._next_delay()
|
| 246 |
+
|
| 247 |
+
def _next_delay(self) -> float:
|
| 248 |
+
"""Poll interval, backed off exponentially while Yahoo is refusing us.
|
| 249 |
+
|
| 250 |
+
Yahoo rate-limits aggressively and an unofficial API gives no
|
| 251 |
+
Retry-After, so a fixed interval just keeps you throttled. Jitter stops
|
| 252 |
+
several symbols (or several deployments) resynchronising after an outage.
|
| 253 |
+
"""
|
| 254 |
+
if self._consecutive_failures == 0:
|
| 255 |
+
base = float(self.poll_interval)
|
| 256 |
+
else:
|
| 257 |
+
base = min(
|
| 258 |
+
self.poll_interval * (2 ** self._consecutive_failures),
|
| 259 |
+
float(self.max_backoff),
|
| 260 |
+
)
|
| 261 |
+
logger.warning(
|
| 262 |
+
"Yahoo poll failed %s time(s) in a row; next attempt in ~%.0fs",
|
| 263 |
+
self._consecutive_failures,
|
| 264 |
+
base,
|
| 265 |
+
)
|
| 266 |
+
return max(1.0, base * random.uniform(0.8, 1.2))
|
| 267 |
+
|
| 268 |
+
def _poll_once(self) -> bool:
|
| 269 |
+
"""Fetch and ingest one round of bars. Returns True if any symbol succeeded."""
|
| 270 |
+
any_success = False
|
| 271 |
for symbol in self.symbols:
|
| 272 |
try:
|
| 273 |
raw = self._download(symbol, period='5d', interval=self.interval)
|
|
|
|
| 276 |
logger.warning("Yahoo poll returned no bars for %s", symbol)
|
| 277 |
continue
|
| 278 |
self._ingest_new_bars(symbol, df)
|
| 279 |
+
any_success = True
|
| 280 |
except Exception as e:
|
| 281 |
logger.error("Yahoo poll failed for %s: %s", symbol, e)
|
| 282 |
+
return any_success
|
| 283 |
+
|
| 284 |
+
def _warn_if_unadjusted(self, symbol: str, df: pd.DataFrame) -> int:
|
| 285 |
+
"""Flag single-bar moves that look like unadjusted corporate actions.
|
| 286 |
+
|
| 287 |
+
This is a backstop rather than the fix -- the fix is auto_adjust. But a
|
| 288 |
+
split slipping through silently corrupts every downstream number, so it
|
| 289 |
+
is worth naming the dates rather than letting a strategy trade them.
|
| 290 |
+
Returns the number of suspicious bars found.
|
| 291 |
+
"""
|
| 292 |
+
duration = _INTERVAL_DURATION.get(self.interval)
|
| 293 |
+
if df.empty or len(df) < 2 or duration is None or duration < pd.Timedelta(days=1):
|
| 294 |
+
return 0
|
| 295 |
+
moves = df['close'].pct_change()
|
| 296 |
+
suspects = df.loc[moves.abs() > _SPLIT_SUSPECT_MOVE, 'timestamp']
|
| 297 |
+
if len(suspects):
|
| 298 |
+
dates = ', '.join(str(pd.Timestamp(t).date()) for t in suspects.head(5))
|
| 299 |
+
logger.warning(
|
| 300 |
+
"%s has %s bar(s) moving more than %.0f%% (%s). On a liquid name that is "
|
| 301 |
+
"usually an unadjusted split, not a real move — check yahoo.auto_adjust.",
|
| 302 |
+
symbol,
|
| 303 |
+
len(suspects),
|
| 304 |
+
_SPLIT_SUSPECT_MOVE * 100,
|
| 305 |
+
dates,
|
| 306 |
+
)
|
| 307 |
+
return int(len(suspects))
|
| 308 |
+
|
| 309 |
+
def _drop_incomplete(self, df: pd.DataFrame) -> pd.DataFrame:
|
| 310 |
+
"""Remove the bar that is still forming.
|
| 311 |
+
|
| 312 |
+
Yahoo returns the in-progress period as an ordinary row. Emitting it
|
| 313 |
+
would hand the strategy a close that has not happened yet, and because
|
| 314 |
+
the watermark advances past it, the finished version never arrives.
|
| 315 |
+
"""
|
| 316 |
+
if self.emit_incomplete_bars or df.empty:
|
| 317 |
+
return df
|
| 318 |
+
duration = _INTERVAL_DURATION.get(self.interval)
|
| 319 |
+
if duration is None:
|
| 320 |
+
return df
|
| 321 |
+
now = pd.Timestamp.now(tz='UTC').tz_convert(None)
|
| 322 |
+
complete = df[df['timestamp'] + duration <= now]
|
| 323 |
+
dropped = len(df) - len(complete)
|
| 324 |
+
if dropped:
|
| 325 |
+
logger.debug("Dropped %s in-progress %s bar(s)", dropped, self.interval)
|
| 326 |
+
return complete
|
| 327 |
|
| 328 |
def _ingest_new_bars(self, symbol: str, df: pd.DataFrame) -> None:
|
| 329 |
+
rows = self._drop_incomplete(df)
|
| 330 |
last_ts = self._last_bar_ts.get(symbol)
|
|
|
|
| 331 |
if last_ts is not None:
|
| 332 |
+
rows = rows[rows['timestamp'] > last_ts]
|
| 333 |
if rows.empty:
|
| 334 |
return
|
| 335 |
|
|
@@ -33,7 +33,13 @@ alpaca:
|
|
| 33 |
# Unofficial API, typically delayed; 1m lookback is ~7 days.
|
| 34 |
yahoo:
|
| 35 |
poll_interval_seconds: 60
|
| 36 |
-
auto_adjust
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 37 |
start_date: '2024-01-01'
|
| 38 |
end_date: '2026-12-31'
|
| 39 |
|
|
|
|
| 33 |
# Unofficial API, typically delayed; 1m lookback is ~7 days.
|
| 34 |
yahoo:
|
| 35 |
poll_interval_seconds: 60
|
| 36 |
+
# Keep true. With auto_adjust off, Yahoo returns raw Close and every stock
|
| 37 |
+
# split reads as a crash (NVDA's June 2024 10:1 becomes a -90% bar).
|
| 38 |
+
auto_adjust: true
|
| 39 |
+
# Yahoo returns the still-forming period as an ordinary row. Emitting it would
|
| 40 |
+
# trade on a close that has not happened yet.
|
| 41 |
+
emit_incomplete_bars: false
|
| 42 |
+
max_backoff_seconds: 900
|
| 43 |
start_date: '2024-01-01'
|
| 44 |
end_date: '2026-12-31'
|
| 45 |
|
|
@@ -1,3 +1,5 @@
|
|
|
|
|
|
|
|
| 1 |
import pandas as pd
|
| 2 |
import pytest
|
| 3 |
from unittest.mock import patch
|
|
@@ -32,6 +34,39 @@ def _sample_yahoo_frame():
|
|
| 32 |
)
|
| 33 |
|
| 34 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 35 |
class TestYahooDataStream:
|
| 36 |
def test_initialization_from_symbol(self, yahoo_config):
|
| 37 |
stream = YahooDataStream(yahoo_config)
|
|
@@ -58,3 +93,137 @@ class TestYahooDataStream:
|
|
| 58 |
df = stream.get_historical_data('AAPL', '2024-01-01', '2024-12-31')
|
| 59 |
assert len(df) == 3
|
| 60 |
assert 'open' in df.columns
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
import logging
|
| 2 |
+
|
| 3 |
import pandas as pd
|
| 4 |
import pytest
|
| 5 |
from unittest.mock import patch
|
|
|
|
| 34 |
)
|
| 35 |
|
| 36 |
|
| 37 |
+
def _split_frame():
|
| 38 |
+
"""An unadjusted 10:1 split, as Yahoo returns it with auto_adjust=False.
|
| 39 |
+
|
| 40 |
+
Modelled on NVDA, 10 June 2024: the raw Close drops from ~1200 to ~120 and
|
| 41 |
+
a backtest reads it as a -90% day.
|
| 42 |
+
"""
|
| 43 |
+
idx = pd.date_range('2024-06-06', periods=4, freq='D', tz='America/New_York')
|
| 44 |
+
close = [1200.0, 1208.0, 120.5, 121.0]
|
| 45 |
+
return pd.DataFrame(
|
| 46 |
+
{
|
| 47 |
+
'Open': close,
|
| 48 |
+
'High': [c * 1.01 for c in close],
|
| 49 |
+
'Low': [c * 0.99 for c in close],
|
| 50 |
+
'Close': close,
|
| 51 |
+
'Volume': [1_000_000] * 4,
|
| 52 |
+
},
|
| 53 |
+
index=idx,
|
| 54 |
+
)
|
| 55 |
+
|
| 56 |
+
|
| 57 |
+
def _dated_frame(timestamps, close=100.0):
|
| 58 |
+
return pd.DataFrame(
|
| 59 |
+
{
|
| 60 |
+
'timestamp': [pd.Timestamp(t) for t in timestamps],
|
| 61 |
+
'open': close,
|
| 62 |
+
'high': close,
|
| 63 |
+
'low': close,
|
| 64 |
+
'close': close,
|
| 65 |
+
'volume': 1_000.0,
|
| 66 |
+
}
|
| 67 |
+
)
|
| 68 |
+
|
| 69 |
+
|
| 70 |
class TestYahooDataStream:
|
| 71 |
def test_initialization_from_symbol(self, yahoo_config):
|
| 72 |
stream = YahooDataStream(yahoo_config)
|
|
|
|
| 93 |
df = stream.get_historical_data('AAPL', '2024-01-01', '2024-12-31')
|
| 94 |
assert len(df) == 3
|
| 95 |
assert 'open' in df.columns
|
| 96 |
+
|
| 97 |
+
|
| 98 |
+
class TestPriceAdjustment:
|
| 99 |
+
"""Unadjusted prices turn every split into a phantom crash."""
|
| 100 |
+
|
| 101 |
+
def test_adjustment_is_on_by_default(self):
|
| 102 |
+
stream = YahooDataStream({'trading': {'symbol': 'AAPL', 'timeframe': '1d'}})
|
| 103 |
+
assert stream.auto_adjust is True
|
| 104 |
+
|
| 105 |
+
def test_auto_adjust_is_passed_through_to_yfinance(self):
|
| 106 |
+
stream = YahooDataStream({'trading': {'symbol': 'AAPL', 'timeframe': '1d'}})
|
| 107 |
+
with patch('yfinance.download', return_value=_sample_yahoo_frame()) as download:
|
| 108 |
+
stream._download('AAPL', period='5d', interval='1d')
|
| 109 |
+
assert download.call_args.kwargs['auto_adjust'] is True
|
| 110 |
+
|
| 111 |
+
def test_opting_out_of_adjustment_warns(self, caplog):
|
| 112 |
+
with caplog.at_level(logging.WARNING):
|
| 113 |
+
YahooDataStream({
|
| 114 |
+
'trading': {'symbol': 'AAPL', 'timeframe': '1d'},
|
| 115 |
+
'yahoo': {'auto_adjust': False},
|
| 116 |
+
})
|
| 117 |
+
assert any('not split' in r.message.lower() for r in caplog.records)
|
| 118 |
+
|
| 119 |
+
def test_split_sized_move_is_flagged(self, yahoo_config, caplog):
|
| 120 |
+
stream = YahooDataStream(yahoo_config)
|
| 121 |
+
df = stream._normalize_ohlcv(_split_frame())
|
| 122 |
+
with caplog.at_level(logging.WARNING):
|
| 123 |
+
found = stream._warn_if_unadjusted('NVDA', df)
|
| 124 |
+
assert found == 1
|
| 125 |
+
assert any('auto_adjust' in r.message for r in caplog.records)
|
| 126 |
+
|
| 127 |
+
def test_ordinary_moves_are_not_flagged(self, yahoo_config, caplog):
|
| 128 |
+
stream = YahooDataStream(yahoo_config)
|
| 129 |
+
df = stream._normalize_ohlcv(_sample_yahoo_frame())
|
| 130 |
+
with caplog.at_level(logging.WARNING):
|
| 131 |
+
assert stream._warn_if_unadjusted('AAPL', df) == 0
|
| 132 |
+
|
| 133 |
+
def test_intraday_bars_are_not_split_checked(self, yahoo_config):
|
| 134 |
+
"""A 40% move in one minute is a halt or a fat finger, not a split."""
|
| 135 |
+
yahoo_config['trading']['timeframe'] = '1m'
|
| 136 |
+
stream = YahooDataStream(yahoo_config)
|
| 137 |
+
df = stream._normalize_ohlcv(_split_frame())
|
| 138 |
+
assert stream._warn_if_unadjusted('NVDA', df) == 0
|
| 139 |
+
|
| 140 |
+
|
| 141 |
+
class TestIncompleteBars:
|
| 142 |
+
"""The bar Yahoo is still building must not be reported as final."""
|
| 143 |
+
|
| 144 |
+
def test_forming_bar_is_dropped(self, yahoo_config):
|
| 145 |
+
stream = YahooDataStream(yahoo_config)
|
| 146 |
+
now = pd.Timestamp.now(tz='UTC').tz_convert(None).normalize()
|
| 147 |
+
df = _dated_frame([now - pd.Timedelta(days=2), now - pd.Timedelta(days=1), now])
|
| 148 |
+
kept = stream._drop_incomplete(df)
|
| 149 |
+
assert len(kept) == 2
|
| 150 |
+
assert kept['timestamp'].max() < now
|
| 151 |
+
|
| 152 |
+
def test_finished_bars_all_survive(self, yahoo_config):
|
| 153 |
+
stream = YahooDataStream(yahoo_config)
|
| 154 |
+
now = pd.Timestamp.now(tz='UTC').tz_convert(None).normalize()
|
| 155 |
+
df = _dated_frame([now - pd.Timedelta(days=5), now - pd.Timedelta(days=4)])
|
| 156 |
+
assert len(stream._drop_incomplete(df)) == 2
|
| 157 |
+
|
| 158 |
+
def test_opting_in_keeps_the_forming_bar(self, yahoo_config):
|
| 159 |
+
yahoo_config['yahoo']['emit_incomplete_bars'] = True
|
| 160 |
+
stream = YahooDataStream(yahoo_config)
|
| 161 |
+
now = pd.Timestamp.now(tz='UTC').tz_convert(None).normalize()
|
| 162 |
+
df = _dated_frame([now - pd.Timedelta(days=1), now])
|
| 163 |
+
assert len(stream._drop_incomplete(df)) == 2
|
| 164 |
+
|
| 165 |
+
def test_partial_bar_is_never_emitted_then_stranded(self, yahoo_config):
|
| 166 |
+
"""The bug this guards: emitting the forming bar advanced the watermark,
|
| 167 |
+
so the finished version of that same bar never reached a callback."""
|
| 168 |
+
stream = YahooDataStream(yahoo_config)
|
| 169 |
+
received = []
|
| 170 |
+
stream.add_data_callback(lambda kind, bar: received.append(bar))
|
| 171 |
+
|
| 172 |
+
now = pd.Timestamp.now(tz='UTC').tz_convert(None).normalize()
|
| 173 |
+
yesterday, today = now - pd.Timedelta(days=1), now
|
| 174 |
+
|
| 175 |
+
stream._ingest_new_bars('AAPL', _dated_frame([yesterday, today], close=100.0))
|
| 176 |
+
assert len(received) == 1 # only yesterday's completed bar
|
| 177 |
+
|
| 178 |
+
# Next day: what was the forming bar is now final and must arrive.
|
| 179 |
+
with patch.object(stream, '_drop_incomplete', side_effect=lambda d: d):
|
| 180 |
+
stream._ingest_new_bars('AAPL', _dated_frame([yesterday, today], close=105.0))
|
| 181 |
+
assert len(received) == 2
|
| 182 |
+
assert received[-1]['close'] == 105.0
|
| 183 |
+
|
| 184 |
+
|
| 185 |
+
class TestPollBackoff:
|
| 186 |
+
"""Yahoo rate-limits hard, and a fixed interval keeps you throttled."""
|
| 187 |
+
|
| 188 |
+
def test_success_polls_at_the_configured_interval(self, yahoo_config):
|
| 189 |
+
yahoo_config['yahoo']['poll_interval_seconds'] = 60
|
| 190 |
+
stream = YahooDataStream(yahoo_config)
|
| 191 |
+
stream._consecutive_failures = 0
|
| 192 |
+
assert 48 <= stream._next_delay() <= 72 # 60s +/- jitter
|
| 193 |
+
|
| 194 |
+
def test_delay_grows_with_consecutive_failures(self, yahoo_config):
|
| 195 |
+
yahoo_config['yahoo']['poll_interval_seconds'] = 60
|
| 196 |
+
stream = YahooDataStream(yahoo_config)
|
| 197 |
+
delays = []
|
| 198 |
+
for failures in (1, 2, 3):
|
| 199 |
+
stream._consecutive_failures = failures
|
| 200 |
+
delays.append(stream._next_delay())
|
| 201 |
+
assert delays[0] < delays[1] < delays[2]
|
| 202 |
+
|
| 203 |
+
def test_backoff_is_capped(self, yahoo_config):
|
| 204 |
+
yahoo_config['yahoo']['poll_interval_seconds'] = 60
|
| 205 |
+
yahoo_config['yahoo']['max_backoff_seconds'] = 300
|
| 206 |
+
stream = YahooDataStream(yahoo_config)
|
| 207 |
+
stream._consecutive_failures = 20
|
| 208 |
+
assert stream._next_delay() <= 300 * 1.2
|
| 209 |
+
|
| 210 |
+
def test_jitter_desynchronises_retries(self, yahoo_config):
|
| 211 |
+
stream = YahooDataStream(yahoo_config)
|
| 212 |
+
stream._consecutive_failures = 3
|
| 213 |
+
assert len({stream._next_delay() for _ in range(20)}) > 1
|
| 214 |
+
|
| 215 |
+
def test_poll_reports_failure_when_every_symbol_fails(self, yahoo_config):
|
| 216 |
+
stream = YahooDataStream(yahoo_config)
|
| 217 |
+
with patch.object(stream, '_download', side_effect=RuntimeError('429 Too Many Requests')):
|
| 218 |
+
assert stream._poll_once() is False
|
| 219 |
+
|
| 220 |
+
def test_poll_reports_success_when_a_symbol_returns_bars(self, yahoo_config):
|
| 221 |
+
stream = YahooDataStream(yahoo_config)
|
| 222 |
+
with patch.object(stream, '_download', return_value=_sample_yahoo_frame()):
|
| 223 |
+
assert stream._poll_once() is True
|
| 224 |
+
|
| 225 |
+
def test_empty_response_counts_as_failure(self, yahoo_config):
|
| 226 |
+
"""A rate-limited yfinance returns an empty frame rather than raising."""
|
| 227 |
+
stream = YahooDataStream(yahoo_config)
|
| 228 |
+
with patch.object(stream, '_download', return_value=pd.DataFrame()):
|
| 229 |
+
assert stream._poll_once() is False
|