@ -19,6 +19,11 @@ from app.schemas.financial import DataSource, ErrorType
from app . utils . date_utils import parse_period , quarters_to_date_range , resolve_time_parameters
from app . utils . date_utils import parse_period , quarters_to_date_range , resolve_time_parameters
from app . core . config import settings
from app . core . config import settings
try :
import pandas as pd
except ImportError :
pd = None # type: ignore
logger = logging . getLogger ( __name__ )
logger = logging . getLogger ( __name__ )
@ -303,6 +308,104 @@ class PriceDataService:
logger . error ( f " Error fetching intraday for { ticker } : { str ( e ) } " )
logger . error ( f " Error fetching intraday for { ticker } : { str ( e ) } " )
raise
raise
async def get_multi_intraday (
self ,
tickers : List [ str ] ,
interval : str = " 5m " ,
start_date : Optional [ date ] = None ,
end_date : Optional [ date ] = None ,
chunk_size : int = 50 ,
) - > Dict [ str , List [ Dict ] ] :
"""
Fetch intraday bars for multiple tickers via yf . download ( ) .
Returns :
Dict mapping ticker → list of { timestamp , open , high , low , close , volume }
"""
if not self . yf_available :
raise ValueError ( " Yahoo Finance data source not available " )
# end date for yfinance download must be exclusive (day after)
from datetime import timedelta
start_str = start_date . isoformat ( ) if start_date else None
end_str = ( end_date + timedelta ( days = 1 ) ) . isoformat ( ) if end_date else None
result : Dict [ str , List [ Dict ] ] = { t . upper ( ) : [ ] for t in tickers }
loop = asyncio . get_event_loop ( )
for i in range ( 0 , len ( tickers ) , chunk_size ) :
chunk = [ t . upper ( ) for t in tickers [ i : i + chunk_size ] ]
_tickers_str = " " . join ( chunk )
try :
bulk_data = await _run_with_timeout (
loop . run_in_executor (
None ,
lambda ts = _tickers_str : yf . download (
tickers = ts ,
start = start_str ,
end = end_str ,
interval = interval ,
auto_adjust = True ,
prepost = False ,
group_by = " ticker " ,
threads = True ,
progress = False ,
) ,
) ,
timeout_seconds = 120 ,
description = f " multi_intraday chunk { i / / chunk_size + 1 } " ,
)
except Exception as e :
logger . error ( f " multi_intraday chunk error: { e } " )
continue
if bulk_data is None or bulk_data . empty :
continue
def _parse_row ( row ) :
def _f ( v ) :
try :
return None if pd . isna ( v ) else float ( v )
except Exception :
return None
return {
" open " : _f ( row . get ( " Open " ) ) ,
" high " : _f ( row . get ( " High " ) ) ,
" low " : _f ( row . get ( " Low " ) ) ,
" close " : _f ( row . get ( " Close " ) ) or 0.0 ,
" volume " : _f ( row . get ( " Volume " ) ) ,
}
if len ( chunk ) == 1 :
# Single-ticker: flat DataFrame
ticker = chunk [ 0 ]
for ts , row in bulk_data . iterrows ( ) :
dt = ts . to_pydatetime ( )
if dt . tzinfo is None :
dt = dt . replace ( tzinfo = timezone . utc )
result [ ticker ] . append ( { " timestamp " : dt . isoformat ( ) , * * _parse_row ( row ) } )
else :
# Multi-ticker: MultiIndex columns grouped by ticker
for ticker in chunk :
try :
lvl0 = bulk_data . columns . get_level_values ( 0 )
if ticker not in lvl0 :
continue
ticker_df = bulk_data [ ticker ]
for ts , row in ticker_df . iterrows ( ) :
dt = ts . to_pydatetime ( )
if dt . tzinfo is None :
dt = dt . replace ( tzinfo = timezone . utc )
result [ ticker ] . append ( { " timestamp " : dt . isoformat ( ) , * * _parse_row ( row ) } )
except Exception as e :
logger . error ( f " multi_intraday parse error for { ticker } : { e } " )
await asyncio . sleep ( 0.1 ) # rate-limit courtesy
return result
async def get_today_ohlc ( self , ticker : str ) - > Dict :
async def get_today_ohlc ( self , ticker : str ) - > Dict :
""" Get today ' s OHLC. If daily not yet finalized, aggregate from intraday 1m. """
""" Get today ' s OHLC. If daily not yet finalized, aggregate from intraday 1m. """
if not self . yf_available :
if not self . yf_available :