Files
alpaca-llm-bot-v1/bot.py

250 lines
9.8 KiB
Python

from apscheduler.schedulers.background import BackgroundScheduler
from datetime import datetime, timedelta
import json
from sqlalchemy import func
from config import settings
from db import SessionLocal, BotDecision, TradeExecution, CuratedInsight
from services import (
searx_news,
extra_research,
summarize_news_with_ollama,
strategy_signals,
llm_final_decision,
place_order,
market_open,
positions_snapshot,
account_snapshot,
qdrant_similar,
qdrant_add_memory,
trilium_log,
n8n_emit,
redis_set_json,
redis_get_json,
redis_lock,
meili_index_doc,
)
scheduler = BackgroundScheduler(timezone=settings.timezone)
def _daily_spent(db):
since = datetime.utcnow() - timedelta(hours=24)
return float(db.query(func.coalesce(func.sum(TradeExecution.notional), 0)).filter(TradeExecution.ts >= since).scalar() or 0)
def curate_cycle():
db = SessionLocal()
try:
for symbol in settings.symbols:
news = searx_news(symbol)
summary = summarize_news_with_ollama(symbol, news)
row = CuratedInsight(symbol=symbol, summary=summary, sources=json.dumps(news)[:60000])
db.add(row)
meili_index_doc('insights', {
'id': f"ins-{int(datetime.utcnow().timestamp())}-{symbol}",
'ts': datetime.utcnow().isoformat(),
'symbol': symbol,
'summary': summary,
'sources': json.dumps(news)[:2000],
'kind': 'insight'
})
db.commit()
finally:
db.close()
def run_cycle():
if not redis_lock('alpaca:run_cycle:lock', ttl=240):
return
db = SessionLocal()
try:
acct = account_snapshot()
# If paper credentials are invalid/missing, avoid flooding failed decisions.
if settings.paper_mode and not acct.get('id'):
trilium_log('Bot paused: paper auth invalid', 'Paper mode enabled but Alpaca paper account auth failed. Skipping run_cycle to avoid false failures.')
return
if not market_open() and not (settings.paper_mode and settings.force_paper_trades):
return
spent = _daily_spent(db)
if spent >= settings.max_daily_notional:
return
pos = positions_snapshot()
if len(pos) >= settings.max_open_positions:
return
for symbol in settings.symbols:
if spent >= settings.max_daily_notional:
break
cached_news = redis_get_json(f'alpaca:news:{symbol}')
news = cached_news if cached_news else searx_news(symbol)
if news:
redis_set_json(f'alpaca:news:{symbol}', {'items': news} if isinstance(news, list) else news, ttl=900)
if isinstance(news, dict) and 'items' in news:
news = news['items']
strat = strategy_signals(symbol, news)
memory_hits = qdrant_similar(symbol, json.dumps(news)[:2000], limit=5)
decision = llm_final_decision(symbol, news, strat, memory_hits)
if decision.get("needs_more_research") or decision.get("confidence", 0) < settings.min_confidence:
more = extra_research(symbol, decision.get("research_topics", []))
if more:
news = news + more
memory_hits = qdrant_similar(symbol, json.dumps(news)[:2000], limit=5)
decision = llm_final_decision(symbol, news, strat, memory_hits)
if settings.paper_mode and settings.force_paper_trades and decision["action"] == "hold":
decision["action"] = "buy"
decision["confidence"] = max(decision.get("confidence", 0.0), settings.min_confidence)
decision["order_usd"] = max(decision.get("order_usd", 0.0), settings.force_paper_min_usd)
decision["reason"] = f"{decision.get('reason','')} | force_paper_trades"
redis_set_json(
f"alpaca:last_decision:{symbol}",
{
"symbol": symbol,
"action": decision["action"],
"confidence": decision["confidence"],
"order_usd": decision["order_usd"],
"reason": decision["reason"],
"ts": datetime.utcnow().isoformat(),
},
ttl=86400,
)
drow = BotDecision(
symbol=symbol,
action=decision["action"],
confidence=decision["confidence"],
reason=f"{decision['reason']} | strat={strat['strategy']} score={strat['score']}",
market_context=json.dumps(news)[:60000],
order_usd=decision["order_usd"],
status="planned",
)
db.add(drow)
db.commit()
db.refresh(drow)
meili_index_doc('decisions', {
'id': f"dec-{drow.id}",
'ts': drow.ts.isoformat() if drow.ts else datetime.utcnow().isoformat(),
'symbol': drow.symbol,
'action': drow.action,
'confidence': drow.confidence,
'status': drow.status,
'reason': drow.reason[:500],
'order_usd': drow.order_usd,
'kind': 'decision'
})
should_trade = decision["action"] in {"buy", "sell"} and decision["confidence"] >= settings.min_confidence
if should_trade:
effective_cost = settings.fee_per_trade_usd + (settings.slippage_bps / 10000.0) * decision["order_usd"]
notional = min(settings.max_order_usd, decision["order_usd"], settings.max_daily_notional - spent)
if notional <= effective_cost:
drow.status = "risk_blocked"
db.add(drow)
db.commit()
continue
res = place_order(symbol, decision["action"], notional)
ok = bool(res and res.get("ok"))
drow.status = "executed" if ok else "failed"
db.add(drow)
if ok:
trade = TradeExecution(
symbol=symbol,
side=decision["action"],
qty=float((res or {}).get("json", {}).get("qty", 0) or 0),
notional=notional,
alpaca_order_id=(res or {}).get("json", {}).get("id", ""),
raw=json.dumps({"decision": decision, "strategy": strat, "memory": memory_hits, "broker": res})[:60000],
)
db.add(trade)
db.flush()
meili_index_doc('trades', {
'id': f"tr-{trade.id}",
'ts': trade.ts.isoformat() if trade.ts else datetime.utcnow().isoformat(),
'symbol': trade.symbol,
'side': trade.side,
'notional': trade.notional,
'order_id': trade.alpaca_order_id,
'kind': 'trade'
})
db.commit()
# learning memory + notes + orchestration signal
qdrant_add_memory(symbol, f"{symbol} {decision['action']} conf={decision['confidence']} reason={decision['reason']}", {
"memory_type": "execution",
"status": drow.status,
"action": decision["action"],
"confidence": decision["confidence"],
"outcome_score": 1 if ok else -1,
"ts": datetime.utcnow().isoformat(),
})
trilium_log(
f"Trade {symbol} {decision['action']} {drow.status}",
f"## Decision\n- symbol: {symbol}\n- action: {decision['action']}\n- confidence: {decision['confidence']:.2f}\n- status: {drow.status}\n- notional: ${notional:.2f}\n\n## Reason\n{decision['reason']}\n\n## Strategy\n{json.dumps(strat, indent=2)}\n"
)
n8n_emit({"event": "trade", "symbol": symbol, "status": drow.status, "decision": decision, "notional": notional})
if ok:
spent += notional
else:
drow.status = "skipped"
db.add(drow)
db.commit()
# store non-trade decisions too for memory
qdrant_add_memory(symbol, f"{symbol} decision={decision['action']} conf={decision['confidence']} status={drow.status}", {
"memory_type": "decision",
"status": drow.status,
"action": decision["action"],
"confidence": decision["confidence"],
"outcome_score": 0 if drow.status in {"skipped", "risk_blocked"} else (-1 if drow.status=="failed" else 1),
"ts": datetime.utcnow().isoformat(),
})
finally:
db.close()
def start_scheduler():
scheduler.add_job(
curate_cycle,
"interval",
minutes=settings.curate_interval_minutes,
id="curate_cycle",
replace_existing=True,
coalesce=True,
max_instances=1,
misfire_grace_time=120,
)
if settings.trade_interval_minutes and settings.trade_interval_minutes > 0:
scheduler.add_job(
run_cycle,
"interval",
minutes=settings.trade_interval_minutes,
id="trade_cycle",
replace_existing=True,
coalesce=True,
max_instances=1,
misfire_grace_time=120,
)
else:
scheduler.add_job(
run_cycle,
"interval",
hours=settings.trade_interval_hours,
id="trade_cycle",
replace_existing=True,
coalesce=True,
max_instances=1,
misfire_grace_time=120,
)
scheduler.start()