diff --git a/backend/app/api/v1/symbols.py b/backend/app/api/v1/symbols.py index ec60aff..7675dfd 100755 --- a/backend/app/api/v1/symbols.py +++ b/backend/app/api/v1/symbols.py @@ -21,11 +21,12 @@ router = APIRouter(prefix="/symbols", tags=["symbols"]) async def list_symbols( exchange: Optional[str] = Query(None, description="Filter by exchange name"), active_only: bool = Query(True, description="Only return active symbols"), + is_trading: Optional[bool] = Query(None, description="Filter by trading status (top-100 bases)"), limit: int = Query(50, ge=1, le=1000, description="Max symbols to return"), offset: int = Query(0, ge=0, description="Offset for pagination"), db: AsyncSession = Depends(get_db_session), ) -> list[SymbolResponse]: - """Return symbols with pagination, optionally filtered by exchange and/or active status.""" + """Return symbols with pagination, optionally filtered by exchange, active status, and trading flag.""" query = select(Symbol).options(joinedload(Symbol.exchange)) if exchange: @@ -34,6 +35,8 @@ async def list_symbols( ) if active_only: query = query.where(Symbol.is_active == True) # noqa: E712 + if is_trading is not None: + query = query.where(Symbol.is_trading == is_trading) query = query.limit(limit).offset(offset) result = await db.execute(query) @@ -47,6 +50,7 @@ async def list_symbols( base=s.base, quote=s.quote, is_active=s.is_active, + is_trading=s.is_trading, ) for s in symbols ] @@ -82,6 +86,7 @@ async def search_symbols( base=s.base, quote=s.quote, is_active=s.is_active, + is_trading=s.is_trading, ) for s in symbols ] diff --git a/backend/app/config.py b/backend/app/config.py index b132725..6ba8fc9 100755 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -18,7 +18,7 @@ class Settings(BaseSettings): # to override whatever password is embedded in DATABASE_URL. Keeps # DB credentials out of plain env vars, consistent with how JWT keys # are already handled. - DB_PASSWORD_FILE: str = "" + DB_PASSWORD_FILE: str = "/run/secrets/db_pw.txt" # JWT JWT_PRIVATE_KEY_PATH: str = "/run/secrets/jwt_private.pem" @@ -30,7 +30,7 @@ class Settings(BaseSettings): # Encryption ENCRYPTION_KEY: str = "" # If set, ENCRYPTION_KEY is read from this file (Docker secret) instead. - ENCRYPTION_KEY_FILE: str = "" + ENCRYPTION_KEY_FILE: str = "/run/secrets/enc_key.txt" # Redis (shared cache across backend-api / backend-scheduler processes). # Optional: if unreachable, callers fall back to per-process in-memory diff --git a/backend/app/schemas/symbol.py b/backend/app/schemas/symbol.py index fca2549..aeb3ae3 100755 --- a/backend/app/schemas/symbol.py +++ b/backend/app/schemas/symbol.py @@ -11,6 +11,7 @@ class SymbolResponse(BaseModel): base: str quote: str is_active: bool = True + is_trading: bool = False class SymbolSearchResponse(BaseModel): diff --git a/backend/scripts/backfill_candles.py b/backend/scripts/backfill_candles.py new file mode 100644 index 0000000..9b6e481 --- /dev/null +++ b/backend/scripts/backfill_candles.py @@ -0,0 +1,156 @@ +#!/usr/bin/env python3 +"""Historical candle backfill — fetch 500 candles for ALL timeframes for is_trading=true symbols.""" +import asyncio +import logging +import sys +from datetime import datetime, timezone + +from sqlalchemy import and_, select +from sqlalchemy.dialects.postgresql import insert as pg_insert +from sqlalchemy.orm import joinedload + +sys.path.insert(0, "/app") + +from app.database import async_session_factory +from app.exchange.factory import factory as exchange_factory +from app.models import Candle, Symbol, Exchange + +logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s") +logger = logging.getLogger("backfill") + +TIMEFRAMES = ["15m", "30m", "1h", "4h", "1d", "1w", "1M"] +FETCH_LIMIT = 500 # Max historical candles per symbol/timeframe +CONCURRENT = 8 # Fetch 8 symbols concurrently +SEM = asyncio.Semaphore(200) # Max concurrent API calls + + +async def fetch_and_store(symbol: Symbol, adapter, sem: asyncio.Semaphore): + """Fetch all timeframes for one symbol and store in DB.""" + exchange_name = symbol.exchange.name + all_values = [] + success_tfs = 0 + + async def _fetch_one_tf(tf: str): + async with sem: + try: + candles = await adapter.fetch_ohlcv( + symbol=symbol.symbol, + timeframe=tf, + limit=FETCH_LIMIT, + ) + return tf, candles + except Exception as e: + err = str(e) + if "BadSymbol" in err or "does not have market symbol" in err: + logger.warning(" %s/%s not found on %s", symbol.symbol, tf, exchange_name) + elif "1M" in err and "invalid" in err.lower(): + logger.debug(" %s: 1M not supported on %s", symbol.symbol, exchange_name) + else: + logger.error(" %s %s/%s: %s", exchange_name, symbol.symbol, tf, e) + return tf, [] + + tasks = [_fetch_one_tf(tf) for tf in TIMEFRAMES] + results = await asyncio.gather(*tasks) + + for tf, candles in results: + if candles: + success_tfs += 1 + for c in candles: + all_values.append({ + "symbol_id": symbol.id, + "timeframe": c.timeframe, + "timestamp": c.timestamp, + "open": c.open, + "high": c.high, + "low": c.low, + "close": c.close, + "volume": c.volume, + }) + + if all_values: + async with async_session_factory() as db: + stmt = pg_insert(Candle).values(all_values) + stmt = stmt.on_conflict_do_nothing( + index_elements=["symbol_id", "timeframe", "timestamp"] + ) + await db.execute(stmt) + await db.commit() + + return success_tfs, len(all_values) + + +async def main(): + # 1. Load all trading symbols + async with async_session_factory() as db: + query = ( + select(Symbol) + .options(joinedload(Symbol.exchange)) + .join(Exchange, Exchange.id == Symbol.exchange_id) + .where( + and_( + Symbol.is_trading == True, + Symbol.is_active == True, + Exchange.is_active == True, + ) + ) + .order_by(Symbol.symbol) + ) + result = await db.execute(query) + symbols = list(result.scalars().all()) + + total_syms = len(symbols) + logger.info("Backfilling %d symbols across 5 exchanges × %d TFs (limit=%d)", + total_syms, len(TIMEFRAMES), FETCH_LIMIT) + logger.info("Estimated: %d × %d × %d = ~%d candles max", + total_syms, len(TIMEFRAMES), FETCH_LIMIT, + total_syms * len(TIMEFRAMES) * FETCH_LIMIT) + + # 2. Group by exchange + by_exchange: dict[str, list] = {} + for sym in symbols: + by_exchange.setdefault(sym.exchange.name, []).append(sym) + + total_candles = 0 + start_time = datetime.now(timezone.utc) + + for ex_name, sym_list in by_exchange.items(): + logger.info("\n=== %s: %d symbols ===", ex_name.upper(), len(sym_list)) + try: + adapter = exchange_factory.create(ex_name) + except ValueError: + logger.error("Unknown exchange: %s", ex_name) + continue + + # Process symbols concurrently (CONCURRENT at a time) + sem = asyncio.Semaphore(CONCURRENT * len(TIMEFRAMES)) + + batch_tasks = [] + for sym in sym_list: + batch_tasks.append(fetch_and_store(sym, adapter, sem)) + + results = await asyncio.gather(*batch_tasks, return_exceptions=True) + + ex_candles = 0 + ex_success = 0 + for i, r in enumerate(results): + if isinstance(r, Exception): + logger.error(" Error: %s", r) + else: + tfs_ok, n = r + if tfs_ok > 0: + ex_success += 1 + ex_candles += n + + total_candles += ex_candles + elapsed = (datetime.now(timezone.utc) - start_time).total_seconds() + logger.info(" %s DONE: %d/%d symbols, %d candles in %.0fs", + ex_name.upper(), ex_success, len(sym_list), ex_candles, elapsed) + + elapsed = (datetime.now(timezone.utc) - start_time).total_seconds() + logger.info("\n=== BACKFILL COMPLETE ===") + logger.info("Total: %d symbols, %d candles, %.0fs (%.1f min)", + total_syms, total_candles, elapsed, elapsed / 60) + + +if __name__ == "__main__": + asyncio.run(main())