Current Signal Context
Historical Score vs BTC Price
Score Bracket Performance
| Score Range | Label | Days | Avg 30d | Avg 90d | Avg 180d | Avg 1yr | Win Rate (1yr) | Max Gain (1yr) | Max Loss (1yr) | Avg Max DD |
|---|
#!/usr/bin/env python3 """ Bitcoin Accumulation Zone Monitor — Web Dashboard FastAPI server with inline HTML/CSS/JS dashboard. Monitors on-chain metrics to identify optimal BTC accumulation zones. """ import asyncio import json import logging import os import sys import threading import time import traceback from contextlib import asynccontextmanager from datetime import datetime, timezone import requests from fastapi import FastAPI from fastapi.responses import HTMLResponse, JSONResponse from pydantic import BaseModel logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(name)s] %(levelname)s: %(message)s") log = logging.getLogger("btc-monitor") BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) sys.path.insert(0, BASE_DIR) from scrapers import fear_greed, price from scoring import engine from dashboard.persistence import ( append_daily_jsonl, atomic_write_json, load_json, load_jsonl_tail, merge_observation, onchain_refresh_due, ) from dashboard.jobs import JobRegistry _shutdown_event = threading.Event() _background_threads = [] _threads_lock = threading.Lock() @asynccontextmanager async def lifespan(_app): """Own background worker startup and graceful shutdown.""" _shutdown_event.clear() scraper_thread = threading.Thread(target=scraper_loop, name="scraper-scheduler") with _threads_lock: _background_threads.append(scraper_thread) scraper_thread.start() try: yield finally: _shutdown_event.set() with _threads_lock: threads = list(_background_threads) for thread in threads: thread.join(timeout=30) with _threads_lock: _background_threads.clear() app = FastAPI(title="Bitcoin Accumulation Zone Monitor", lifespan=lifespan) CONFIG_DIR = os.path.join(BASE_DIR, "config") DATA_DIR = os.path.join(BASE_DIR, "data") CACHE_PATH = os.path.join(DATA_DIR, "cache.json") HISTORY_PATH = os.path.join(DATA_DIR, "score_history.jsonl") LLM_SETTINGS_PATH = os.path.join(CONFIG_DIR, "llm_settings.json") JOBS_PATH = os.path.join(DATA_DIR, "jobs.json") os.makedirs(DATA_DIR, exist_ok=True) _jobs = JobRegistry(JOBS_PATH) # Background scraper state _scraper_lock = threading.Lock() _scraper_running = False _last_update = None _last_error = None def _job_worker(job_id, operation): try: _jobs.run(job_id, operation) except Exception: log.error("Background job %s failed:\n%s", job_id, traceback.format_exc()) finally: current = threading.current_thread() with _threads_lock: if current in _background_threads: _background_threads.remove(current) def _spawn_job(job, operation): """Start an already-reserved job in a tracked, non-daemon thread.""" thread = threading.Thread( target=_job_worker, args=(job["id"], operation), name=f"{job['kind']}-{job['id'][:8]}", ) with _threads_lock: _background_threads.append(thread) thread.start() return thread # ── Cache management ────────────────────────────────────────────────────── def load_cache(): return load_json(CACHE_PATH, {}) def save_cache(data): atomic_write_json(CACHE_PATH, data) @app.get("/health/live") def health_live(): """Report that the API process is responsive.""" return {"status": "ok"} @app.get("/health/ready") def health_ready(): """Report readiness only after a usable score has been persisted.""" scored = load_cache().get("_scored", {}) score = scored.get("composite_score") count = scored.get("scored_count", 0) if score is None or count < 1: return JSONResponse( {"status": "not_ready", "reason": "no usable persisted score"}, status_code=503, ) return {"status": "ready", "score": score, "scored_metrics": count} def append_history(score_data): """Append a daily score entry to history.""" entry = { "timestamp": datetime.now(timezone.utc).isoformat(), "composite_score": score_data.get("composite_score", 0), "scored_count": score_data.get("scored_count", 0), "metrics": { m["key"]: {"score": m["score"], "value": m["value"]} for m in score_data.get("metrics", []) }, } append_daily_jsonl(HISTORY_PATH, entry) def load_history(): return load_jsonl_tail(HISTORY_PATH, limit=90) # ── Background scraper ──────────────────────────────────────────────────── def _scrape_onchain_sources(): """Run independent on-chain providers so one outage cannot mask the other.""" observations = {} errors = [] successful_sources = 0 providers = ( ("LookIntoBitcoin", "scrapers.lookintobitcoin"), ("CheckOnChain", "scrapers.checkonchain"), ) for display_name, module_name in providers: try: module = __import__(module_name, fromlist=["scrape_all"]) observations.update(module.scrape_all()) successful_sources += 1 except Exception as exc: log.error("%s scraping failed: %s\n%s", display_name, exc, traceback.format_exc()) errors.append(f"{display_name}: {exc}") return observations, errors, successful_sources def run_scrape(force_full=False): """Run a scrape cycle and update cache. By default, only refreshes fast data (price, F&G) and reuses cached on-chain data. On-chain metrics (Playwright scrapes) only refresh if: - force_full=True (manual full refresh) - No cached on-chain data exists - Cached on-chain data is >6 hours old (they update daily) """ global _last_update, _last_error, _scraper_running with _scraper_lock: if _scraper_running: return _scraper_running = True try: existing_cache = load_cache() metrics = {} cycle_errors = [] # Fast metrics fail independently so partial outages retain last-known-good data. log.info("Fetching Fear & Greed...") try: metrics["fear_greed"] = merge_observation( existing_cache.get("fear_greed"), fear_greed.fetch(), source="alternative.me" ) except Exception as e: cycle_errors.append(f"Fear & Greed: {e}") metrics["fear_greed"] = merge_observation( existing_cache.get("fear_greed"), None, source="alternative.me", error=str(e), ) log.info("Fetching BTC price...") try: price_current = price.fetch_current() metrics["price"] = merge_observation( existing_cache.get("price"), price_current, source="coingecko" ) except Exception as e: cycle_errors.append(f"Price: {e}") metrics["price"] = merge_observation( existing_cache.get("price"), None, source="coingecko", error=str(e) ) price_current = metrics["price"] log.info("Fetching BTC ATH...") try: ath_data = price.fetch_ath() except Exception as e: cycle_errors.append(f"ATH: {e}") ath_data = {} ath_val = ath_data.get("ath") or existing_cache.get("drawdown", {}).get("ath") if price_current.get("price") and ath_val: drawdown = price.calculate_drawdown(price_current["price"], ath_val) metrics["drawdown"] = {"value": drawdown, "ath": ath_val} elif existing_cache.get("drawdown", {}).get("value") is not None: log.info("ATH fetch failed — reusing cached drawdown") metrics["drawdown"] = existing_cache["drawdown"] else: metrics["drawdown"] = {"value": None} log.info("Fetching historical prices for 200D SMA / Mayer...") try: hist = price.fetch_historical() except Exception as e: cycle_errors.append(f"Historical price: {e}") hist = [] if hist: sma_200d = price.calculate_200d_sma(hist) mayer = price.calculate_mayer_multiple(price_current.get("price"), sma_200d) metrics["price_extras"] = {"sma_200d": sma_200d, "mayer_multiple": mayer} else: # CoinGecko rate-limited — compute from history.json instead try: hist_path = os.path.join(DATA_DIR, "history.json") with open(hist_path) as f: hdata = json.load(f) btc_vals = hdata.get("btc_price", {}).get("values", []) if len(btc_vals) >= 200: sma_200d = sum(btc_vals[-200:]) / 200 cur_p = price_current.get("price") or btc_vals[-1] mayer = cur_p / sma_200d if sma_200d else None metrics["price_extras"] = {"sma_200d": sma_200d, "mayer_multiple": round(mayer, 4) if mayer else None} log.info("Computed 200D SMA from history.json (CoinGecko rate-limited)") elif existing_cache.get("price_extras"): metrics["price_extras"] = existing_cache["price_extras"] except Exception: if existing_cache.get("price_extras"): metrics["price_extras"] = existing_cache["price_extras"] log.info("Reusing cached price_extras") # 3. On-chain metrics — use cached values (historical data is permanent) onchain_keys = ["puell_multiple", "mvrv_zscore", "reserve_risk", "rhodl_ratio", "nupl", "200w_sma", "lth_realized_price", "hash_ribbons", "pi_cycle_bottom", "lth_supply", "sopr", "sellside_risk", "active_address_momentum", "txcount_momentum", "nvt_price", "vdd_multiple"] refresh_onchain = force_full or onchain_refresh_due(existing_cache.get("_onchain_timestamp")) if refresh_onchain: log.info("Refreshing on-chain metrics (forced, missing, or TTL expired)...") onchain, onchain_errors, successful_sources = _scrape_onchain_sources() cycle_errors.extend(onchain_errors) checkonchain_keys = {"sopr", "sellside_risk", "active_address_momentum", "txcount_momentum", "nvt_price", "vdd_multiple"} for key in onchain_keys: source = "checkonchain" if key in checkonchain_keys else "lookintobitcoin" metrics[key] = merge_observation( existing_cache.get(key), onchain.get(key), source=source, error="metric missing from scrape", ) if successful_sources: metrics["_onchain_timestamp"] = datetime.now(timezone.utc).isoformat() elif "_onchain_timestamp" in existing_cache: metrics["_onchain_timestamp"] = existing_cache["_onchain_timestamp"] else: # Reuse cached on-chain values — they're stored permanently log.info("Reusing cached on-chain data (use Full Refresh to re-scrape)") for k in onchain_keys: if k in existing_cache: metrics[k] = existing_cache[k] if "_onchain_timestamp" in existing_cache: metrics["_onchain_timestamp"] = existing_cache["_onchain_timestamp"] # 4. Score everything (classic + ML) log.info("Scoring metrics...") scored = engine.score_all(metrics) metrics["_scored"] = scored # ML-optimized scoring (parallel) try: scored_ml = engine.score_all_ml(metrics) metrics["_scored_ml"] = scored_ml except Exception as e: log.warning("ML scoring failed (non-critical): %s", e) metrics["_timestamp"] = datetime.now(timezone.utc).isoformat() save_cache(metrics) append_history(scored) # Append today's values to permanent history (incremental, not full re-scrape) try: from scrapers.history_updater import update_history update_history() except Exception as e: log.warning("History update failed (non-critical): %s", e) _last_update = datetime.now(timezone.utc).isoformat() _last_error = "; ".join(cycle_errors) if cycle_errors else None log.info("Scrape cycle complete. Composite score: %s", scored["composite_score"]) except Exception as e: log.error("Scrape cycle error: %s\n%s", e, traceback.format_exc()) _last_error = str(e) finally: with _scraper_lock: _scraper_running = False def _run_scheduled_refresh(force_full=False): job = _jobs.reserve("refresh", details={"full": force_full, "scheduled": True}) if job is not None: _jobs.run(job["id"], lambda: run_scrape(force_full=force_full)) def scraper_loop(): """Background loop: refresh quickly every 15 minutes, with on-chain TTL handling.""" cache = load_cache() has_data = any(cache.get(k, {}).get("value") is not None for k in ["puell_multiple", "mvrv_zscore", "nupl"]) _run_scheduled_refresh(force_full=not has_data) while not _shutdown_event.wait(900): _run_scheduled_refresh() # ── LLM Settings (preserved from original) ─────────────────────────────── class LLMSettingsUpdate(BaseModel): provider: str model: str providers: dict class TestConnectionRequest(BaseModel): provider: str providers: dict class FetchModelsRequest(BaseModel): provider: str providers: dict def _load_llm_settings(): if os.path.exists(LLM_SETTINGS_PATH): with open(LLM_SETTINGS_PATH) as f: return json.load(f) return { "provider": "ollama", "model": "qwen3.5:27b", "providers": { "ollama": {"base_url": "http://100.100.242.21:11434"}, "lmstudio": {"base_url": "http://100.100.242.21:1234"}, "openai": {"api_key": ""}, "anthropic": {"api_key": ""}, "openrouter": {"api_key": ""}, }, } def _mask_api_key(key): if not key or len(key) < 8: return "" return "••••••••" + key[-4:] def _safe_settings(settings): out = json.loads(json.dumps(settings)) for name, cfg in out.get("providers", {}).items(): if "api_key" in cfg: cfg["api_key"] = _mask_api_key(cfg["api_key"]) return out def _merge_api_keys(new_providers, existing_providers): for name, cfg in new_providers.items(): if "api_key" in cfg: masked = cfg["api_key"] if masked.startswith("••••") or masked == "": existing_key = existing_providers.get(name, {}).get("api_key", "") cfg["api_key"] = existing_key def _fetch_models(provider, providers): cfg = providers.get(provider, {}) if provider == "ollama": base_url = cfg.get("base_url", "http://100.100.242.21:11434") resp = requests.get(f"{base_url}/api/tags", timeout=10) resp.raise_for_status() return [{"id": m["name"], "name": m["name"]} for m in resp.json().get("models", [])] elif provider == "lmstudio": base_url = cfg.get("base_url", "http://100.100.242.21:1234") resp = requests.get(f"{base_url}/v1/models", timeout=10) resp.raise_for_status() return [{"id": m["id"], "name": m["id"]} for m in resp.json().get("data", [])] elif provider == "openai": api_key = cfg.get("api_key", "") if not api_key: raise ValueError("OpenAI API key is required") resp = requests.get("https://api.openai.com/v1/models", headers={"Authorization": f"Bearer {api_key}"}, timeout=15) resp.raise_for_status() models = [m for m in resp.json().get("data", []) if m["id"].startswith("gpt-")] models.sort(key=lambda m: m["id"]) return [{"id": m["id"], "name": m["id"]} for m in models] elif provider == "anthropic": api_key = cfg.get("api_key", "") if not api_key: raise ValueError("Anthropic API key is required") resp = requests.get("https://api.anthropic.com/v1/models", headers={"x-api-key": api_key, "anthropic-version": "2023-06-01"}, timeout=15) resp.raise_for_status() return [{"id": m["id"], "name": m.get("display_name", m["id"])} for m in resp.json().get("data", [])] elif provider == "openrouter": resp = requests.get("https://openrouter.ai/api/v1/models", timeout=15) resp.raise_for_status() models = resp.json().get("data", []) models.sort(key=lambda m: m.get("id", "")) return [{"id": m["id"], "name": m.get("name", m["id"])} for m in models[:200]] else: raise ValueError(f"Unknown provider: {provider}") # ── API Routes ──────────────────────────────────────────────────────────── def _with_informational_onchain_metrics(scored, cache): """Add non-scored on-chain data cards without changing composite scoring.""" if not isinstance(scored, dict): return scored enriched = dict(scored) metrics = [dict(m) for m in scored.get("metrics", [])] existing_keys = {m.get("key") for m in metrics} lth_supply = cache.get("lth_supply", {}) lth_value = lth_supply.get("value") if lth_value is not None and "lth_supply" not in existing_keys: trend = lth_supply.get("trend") trend_text = f" — {trend}" if trend else "" metrics.append({ "name": "Long-Term Holder Supply", "key": "lth_supply", "value": lth_value, "display_value": f"{lth_value:,.0f} BTC", "score": None, "description": "Informational on-chain metric; not included in the composite score" + trend_text, "recent": lth_supply.get("recent", []), }) pi_cycle = cache.get("pi_cycle_bottom", {}) pi_value = pi_cycle.get("value") if pi_value is not None and "pi_cycle_bottom" not in existing_keys: metrics.append({ "name": "Pi Cycle Bottom", "key": "pi_cycle_bottom", "value": pi_value, "display_value": f"{pi_value:,.2f}" if isinstance(pi_value, (int, float)) else str(pi_value), "score": None, "description": "Informational on-chain cycle metric; not included in the composite score", "recent": pi_cycle.get("recent", []), }) enriched["metrics"] = metrics return enriched @app.get("/api/data") def api_data(mode: str = "classic"): """Return current cached metrics + scores. mode=classic (default) or mode=ml for ML-optimized scoring. """ cache = load_cache() if mode == "ml": scored = cache.get("_scored_ml", cache.get("_scored", {})) else: scored = cache.get("_scored", {}) scored = _with_informational_onchain_metrics(scored, cache) price_data = cache.get("price", {}) drawdown_data = cache.get("drawdown", {}) extras = cache.get("price_extras", {}) return { "scored": scored, "price": price_data.get("price"), "change_24h": price_data.get("change_24h"), "ath": drawdown_data.get("ath"), "mayer_multiple": extras.get("mayer_multiple"), "sma_200d": extras.get("sma_200d"), "last_update": cache.get("_timestamp"), "scraper_running": _scraper_running, "last_error": _last_error, "mode": mode, } @app.get("/api/history") def api_history(): return load_history()[-90:] # Last 90 entries @app.post("/api/refresh", status_code=202) def api_refresh(full: bool = False): """Atomically reserve and start a quick or full metric refresh.""" job = _jobs.reserve("refresh", details={"full": full, "scheduled": False}) if job is None: active = _jobs.active("refresh") return JSONResponse( {"error": "Scrape already in progress", "job": active}, status_code=409 ) _spawn_job(job, lambda: run_scrape(force_full=full)) mode = "full (on-chain + price + F&G)" if full else "quick (price + F&G only)" return { "ok": True, "job_id": job["id"], "status": job["status"], "message": f"Scrape started — {mode}", } @app.get("/api/jobs/{job_id}") def api_job_status(job_id: str): job = _jobs.get(job_id) if job is None: return JSONResponse({"error": "Job not found"}, status_code=404) return job # Settings routes (preserved) @app.get("/api/settings") def api_get_settings(): return _safe_settings(_load_llm_settings()) @app.post("/api/settings") def api_save_settings(body: LLMSettingsUpdate): existing = _load_llm_settings() new_settings = {"provider": body.provider, "model": body.model, "providers": body.providers} _merge_api_keys(new_settings["providers"], existing.get("providers", {})) with open(LLM_SETTINGS_PATH, "w") as f: json.dump(new_settings, f, indent=2) return {"ok": True, "message": "Settings saved"} @app.post("/api/settings/test") def api_test_connection(body: TestConnectionRequest): existing = _load_llm_settings() providers = json.loads(json.dumps(body.providers)) _merge_api_keys(providers, existing.get("providers", {})) try: models = _fetch_models(body.provider, providers) return {"ok": True, "models": models, "message": f"Connected — {len(models)} model(s) found"} except requests.exceptions.ConnectionError: return JSONResponse({"ok": False, "error": "Connection refused"}, status_code=502) except Exception as e: return JSONResponse({"ok": False, "error": str(e)}, status_code=500) @app.post("/api/settings/models") def api_fetch_models(body: FetchModelsRequest): existing = _load_llm_settings() providers = json.loads(json.dumps(body.providers)) _merge_api_keys(providers, existing.get("providers", {})) try: models = _fetch_models(body.provider, providers) return {"ok": True, "models": models} except Exception as e: return JSONResponse({"ok": False, "error": str(e)}, status_code=500) # ── HTML Pages ──────────────────────────────────────────────────────────── SHARED_CSS = """ *,*::before,*::after{box-sizing:border-box;margin:0;padding:0} :root{--bg:#0f172a;--card:#1e293b;--card-hover:#253349;--text:#e2e8f0;--text-dim:#94a3b8; --accent:#f7931a;--green:#22c55e;--red:#ef4444;--yellow:#eab308;--border:#334155; --mono:'JetBrains Mono','Fira Code','Courier New',monospace;--cyan:#22d3ee; --bright-green:#4ade80;--score-excellent:#22c55e;--score-good:#4ade80; --score-neutral:#eab308;--score-bad:#f97316;--score-terrible:#ef4444} body{font-family:'Inter',sans-serif;background:var(--bg);color:var(--text);min-height:100vh} .container{max-width:1400px;margin:0 auto;padding:16px} h1{font-size:1.5rem;font-weight:700;display:flex;align-items:center;gap:10px} h1 .btc{color:var(--accent);font-size:1.8rem} h2{font-size:.8rem;font-weight:600;color:var(--text-dim);margin-bottom:12px;text-transform:uppercase;letter-spacing:.05em} .header{display:flex;justify-content:space-between;align-items:center;padding:16px 0;border-bottom:1px solid var(--border);margin-bottom:16px;flex-wrap:wrap;gap:12px} .nav{display:flex;gap:4px;align-items:center} .nav a{color:var(--text-dim);text-decoration:none;font-size:.85rem;font-weight:600;padding:6px 14px;border-radius:6px;transition:all .15s} .nav a:hover{color:var(--text);background:var(--card)} .nav a.active{color:var(--cyan);background:var(--card);border:1px solid var(--border)} .btn{padding:8px 18px;border:none;border-radius:6px;font-family:inherit;font-weight:600;font-size:.85rem;cursor:pointer;transition:all .15s} .btn-accent{background:var(--accent);color:#000}.btn-accent:hover{background:#e8850f} .btn-secondary{background:var(--border);color:var(--text)}.btn-secondary:hover{background:var(--card-hover)} .btn-cyan{background:var(--cyan);color:#000}.btn-cyan:hover{background:#06b6d4} .btn:disabled{opacity:.4;cursor:not-allowed} .card{background:var(--card);border-radius:10px;padding:16px;border:1px solid var(--border)} .footer{text-align:center;color:var(--text-dim);font-size:.75rem;padding:20px 0;margin-top:16px;border-top:1px solid var(--border)} .toast{position:fixed;top:20px;right:20px;padding:12px 20px;border-radius:8px;font-size:.85rem;font-weight:600;z-index:9999;opacity:0;transform:translateY(-10px);transition:all .3s;pointer-events:none} .toast.show{opacity:1;transform:translateY(0)} .toast-success{background:var(--green);color:#000} .toast-error{background:var(--red);color:#fff} """ SHARED_HEAD = """ """ NAV_HTML = """
""" TOAST_JS = """ function showToast(msg, type) { let t = document.getElementById('toast'); if (!t) { t = document.createElement('div'); t.id = 'toast'; t.className = 'toast'; document.body.appendChild(t); } t.textContent = msg; t.className = 'toast toast-' + type + ' show'; setTimeout(() => t.classList.remove('show'), 3500); } """ DASHBOARD_HTML = """ """ + SHARED_HEAD + """