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