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

200 lines
8.1 KiB
Python

import json
import random
import time
import uuid
import requests
from config import settings
def _ollama_generate(model: str, payload_obj: dict, timeout: int = 45):
payload = {"model": model, "stream": False, "prompt": json.dumps(payload_obj), "format": "json"}
r = requests.post(f"{settings.ollama_url}/api/generate", json=payload, timeout=timeout)
r.raise_for_status()
return json.loads(r.json().get("response", "{}"))
def _embed(text: str):
try:
r = requests.post(f"{settings.ollama_url}/api/embeddings", json={"model": settings.ollama_embed_model, "prompt": text[:8000]}, timeout=30)
r.raise_for_status()
return r.json().get("embedding", [])
except Exception:
return []
def qdrant_ensure_collection(vector_size=768):
try:
requests.put(f"{settings.qdrant_url}/collections/{settings.qdrant_collection}", json={"vectors": {"size": vector_size, "distance": "Cosine"}}, timeout=10)
except Exception:
pass
def qdrant_add_memory(symbol: str, text: str, payload: dict):
vec = _embed(text)
if not vec:
return False
qdrant_ensure_collection(len(vec))
body = {
"points": [{"id": str(uuid.uuid4()), "vector": vec, "payload": {"symbol": symbol, "text": text[:2000], **payload}}]
}
try:
r = requests.put(f"{settings.qdrant_url}/collections/{settings.qdrant_collection}/points", json=body, timeout=15)
return r.ok
except Exception:
return False
def qdrant_similar(symbol: str, query_text: str, limit: int = 5):
vec = _embed(query_text)
if not vec:
return []
body = {"vector": vec, "limit": limit, "with_payload": True}
try:
r = requests.post(f"{settings.qdrant_url}/collections/{settings.qdrant_collection}/points/search", json=body, timeout=15)
if not r.ok:
return []
out = []
for p in r.json().get("result", []):
pl = p.get("payload", {})
if pl.get("symbol") in (symbol, None):
out.append({"score": p.get("score", 0), "text": pl.get("text", ""), "status": pl.get("status", "")})
return out
except Exception:
return []
def searx_news(symbol: str, limit: int = 12):
q = f"{symbol} stock news earnings guidance analyst macro risk"
params = {"q": q, "format": "json", "language": "en"}
try:
r = requests.get(settings.searx_url, params=params, timeout=20)
r.raise_for_status()
data = r.json()
return [{"title": it.get("title", ""), "url": it.get("url", ""), "content": (it.get("content", "") or "")[:700]} for it in data.get("results", [])[:limit]]
except Exception:
return []
def extra_research(symbol: str, weak_points: list, limit: int = 6):
q = f"{symbol} {' '.join(weak_points[:3])} SEC filing guidance risks competition"
params = {"q": q, "format": "json", "language": "en"}
try:
r = requests.get(settings.searx_url, params=params, timeout=20)
r.raise_for_status()
data = r.json()
return [{"title": it.get("title", ""), "url": it.get("url", ""), "content": (it.get("content", "") or "")[:700]} for it in data.get("results", [])[:limit]]
except Exception:
return []
def summarize_news_with_ollama(symbol: str, context_items: list):
prompt = {"task": "Summarize market-moving info into a concise brief", "symbol": symbol, "news": context_items}
try:
parsed = _ollama_generate(settings.ollama_curator_model, prompt)
return parsed.get("summary", "no-summary")
except Exception:
return "summary-unavailable"
def strategy_signals(symbol: str, context_items: list):
text_blob = " ".join((x.get("title", "") + " " + x.get("content", "")) for x in context_items).lower()
bullish = sum(k in text_blob for k in ["beat", "raise guidance", "upgrade", "buyback", "record revenue"])
bearish = sum(k in text_blob for k in ["miss", "downgrade", "lawsuit", "probe", "cut guidance", "recall"])
score = bullish - bearish
action = "buy" if score >= 2 else ("sell" if score <= -2 else "hold")
conf = min(0.85, 0.50 + abs(score) * 0.08)
return {"strategy": "event-momentum-v1", "score": score, "action": action, "confidence": conf, "signals": {"bullish": bullish, "bearish": bearish}}
def llm_final_decision(symbol: str, context_items: list, strategy: dict, memory_hits: list):
prompt = {
"task": "Final trading decision. Return strict JSON.",
"symbol": symbol,
"constraints": {"actions": ["buy", "sell", "hold"], "max_order_usd": settings.max_order_usd, "fee_per_trade_usd": settings.fee_per_trade_usd, "slippage_bps": settings.slippage_bps, "avoid_overtrading": True},
"strategy_prior": strategy,
"memory_hits": memory_hits,
"news": context_items,
"output_schema": {"action": "buy|sell|hold", "confidence": "0-1", "order_usd": f"<= {settings.max_order_usd}", "reason": "short rationale", "needs_more_research": True, "research_topics": ["..."]},
}
try:
d = _ollama_generate(settings.ollama_decision_model, prompt, timeout=60)
action = str(d.get("action", "hold")).lower()
if action not in {"buy", "sell", "hold"}:
action = "hold"
return {
"action": action,
"confidence": max(0.0, min(1.0, float(d.get("confidence", 0.5)))),
"order_usd": min(float(d.get("order_usd", settings.max_order_usd)), settings.max_order_usd),
"reason": d.get("reason", "fallback"),
"needs_more_research": bool(d.get("needs_more_research", False)),
"research_topics": d.get("research_topics", []) or [],
}
except Exception:
return {"action": strategy.get("action", "hold"), "confidence": min(strategy.get("confidence", 0.5), 0.55), "order_usd": min(1.0, settings.max_order_usd), "reason": "decision-fallback-strategy", "needs_more_research": False, "research_topics": []}
def alpaca_headers():
return {"APCA-API-KEY-ID": settings.alpaca_key, "APCA-API-SECRET-KEY": settings.alpaca_secret, "Content-Type": "application/json"}
def place_order(symbol: str, action: str, order_usd: float):
if action not in {"buy", "sell"}:
return None
payload = {"symbol": symbol, "side": action, "type": "market", "time_in_force": "day", "notional": round(order_usd, 2)}
try:
r = requests.post(f"{settings.alpaca_base}/v2/orders", headers=alpaca_headers(), json=payload, timeout=20)
return {"ok": r.ok, "status": r.status_code, "json": r.json() if r.text else {}}
except Exception as e:
return {"ok": False, "status": 0, "json": {"error": str(e)}}
def account_snapshot():
try:
r = requests.get(f"{settings.alpaca_base}/v2/account", headers=alpaca_headers(), timeout=20)
r.raise_for_status()
return r.json()
except Exception:
return {}
def positions_snapshot():
try:
r = requests.get(f"{settings.alpaca_base}/v2/positions", headers=alpaca_headers(), timeout=20)
return r.json() if r.ok else []
except Exception:
return []
def market_open():
try:
r = requests.get(f"{settings.alpaca_base}/v2/clock", headers=alpaca_headers(), timeout=20)
return bool(r.json().get("is_open", False)) if r.ok else False
except Exception:
return False
def trilium_log(title: str, body: str):
if not settings.trilium_token:
return False
headers = {"Authorization": settings.trilium_token, "Content-Type": "application/json"}
payload = {"title": title[:120], "type": "text", "mime": "text/markdown", "content": body}
try:
# best-effort endpoints across Trilium variants
for ep in ["/etapi/create-note", "/etapi/notes"]:
r = requests.post(settings.trilium_url.rstrip("/") + ep, headers=headers, json=payload, timeout=15)
if r.ok:
return True
except Exception:
pass
return False
def n8n_emit(event: dict):
if not settings.n8n_webhook:
return False
try:
r = requests.post(settings.n8n_webhook, json=event, timeout=10)
return r.ok
except Exception:
return False