osint-dashboard/news/summerizer/summarizer.py
Sirius DevOps 1c47ecbc5a feat: load k8s news prompts from env/files, not Python
Bundled MAP_PROMPT/SUMMARY_PROMPT from the customer1 deepseek
configmap. Env wins over prompt files; blank compose injection
is treated as unset. Parser maps market_overview JSON onto the
dashboard ticker/map contract.
2026-08-28 20:04:47 -04:00

432 lines
16 KiB
Python

#!/usr/bin/env python3
"""News summarizer — Nous map-reduce of scraped articles into brief/ticker/map.
Reads articles scraped within the last hour from the shared `articles` table,
maps them with Nous (per-article English fact blocks), reduces to one JSON
object (summary_en + ticker + map_items), and stores the brief in
`article_summaries` plus flagged rows in `news_items`. Tables live in the
EXISTING osint-db (alembic 003_news + 005_news_items, idempotent).
Everything is env-driven (12-factor). Secrets/config are resolved at the start
of each summarize_news() — env wins, else api_keys / app_settings:
DB_HOST / DB_NAME / DB_USER / DB_PASSWORD / DB_PORT PostgreSQL (osint-db)
NOUS_API_KEY Nous Portal key (else api_keys.name='NOUS_API_KEY')
NOUS_BASE_URL default https://inference-api.nousresearch.com/v1
SUMMARY_MODEL default Hermes-4.3-36B (else app_settings)
BATCH_SIZE articles per map-phase batch (default 50)
SUMMARY_WINDOW_HOURS look-back window in hours (default 1)
OSINT_USER_AGENT default osint-dashboard-news-summarizer
MAP_PROMPT override map-phase prompt (uses {batch_text})
SUMMARY_PROMPT override reduce-phase prompt (uses {final_input})
MAP_PROMPT_FILE path to map prompt (default /app/prompt_files/map.txt)
SUMMARY_PROMPT_FILE path to reduce prompt
NEWS_SUMMARIZE_FORCE "1" to ignore the current-UTC-hour idempotency skip
INCLUDE_FUTURES "1" to prepend live futures prices (default 0)
The futures/markets coupling from the original pipeline is gated behind
INCLUDE_FUTURES and OFF by default — it is irrelevant to the OSINT dashboard
and pulled yfinance into the image. Re-enable by installing yfinance and
setting INCLUDE_FUTURES=1.
"""
from __future__ import annotations
import logging
import os
from datetime import datetime
import psycopg2
from intel import parse_reduce_json, select_map, select_ticker
from nous_client import chat
from prompts import map_prompt_template, summary_prompt_template
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
logger = logging.getLogger("news.summarizer")
# ── Configuration (12-factor, container-friendly defaults) ─────────────────
DB_CONFIG = {
"host": os.getenv("DB_HOST", "db").strip(),
"database": os.getenv("DB_NAME", "osint_data").strip(),
"user": os.getenv("DB_USER", "osint").strip(),
"password": os.getenv("DB_PASSWORD", "").strip(),
"port": int(os.getenv("DB_PORT", "5432")),
}
DEFAULT_NOUS_BASE_URL = "https://inference-api.nousresearch.com/v1"
DEFAULT_SUMMARY_MODEL = "Hermes-4.3-36B"
BATCH_SIZE = int(os.getenv("BATCH_SIZE", "50"))
SUMMARY_WINDOW_HOURS = int(os.getenv("SUMMARY_WINDOW_HOURS", "1"))
INCLUDE_FUTURES = os.getenv("INCLUDE_FUTURES", "0").lower() in ("1", "true", "yes")
# Only touched when INCLUDE_FUTURES=1 (legacy markets coupling, OSINT-off).
FUTURES_TICKERS = {
"Equity Indices": ["ES=F", "NQ=F", "YM=F", "RTY=F"],
"Energy": ["CL=F", "NG=F", "HO=F", "RB=F"],
"Metals": ["GC=F", "SI=F", "HG=F"],
"Agriculture": ["ZC=F", "ZS=F", "ZW=F", "ZL=F", "KE=F"],
"Currencies": ["6E=F", "6J=F", "6B=F"],
}
# Prompts live in prompt_files/ (k8s deepseek-configmap). Override at runtime
# with MAP_PROMPT / SUMMARY_PROMPT (env wins) or MAP_PROMPT_FILE / SUMMARY_PROMPT_FILE.
# ── LLM helpers ────────────────────────────────────────────────────────────
def _kv(conn, table, name) -> str:
cur = conn.cursor()
cur.execute(f"SELECT value FROM {table} WHERE name = %s", (name,))
row = cur.fetchone()
return (row[0] or "").strip() if row else ""
def resolve_api_key() -> str:
env = os.getenv("NOUS_API_KEY", "").strip()
if env:
return env
try:
conn = psycopg2.connect(**DB_CONFIG)
try:
return _kv(conn, "api_keys", "NOUS_API_KEY")
finally:
conn.close()
except Exception: # noqa: BLE001
return ""
def resolve_model() -> str:
env = os.getenv("SUMMARY_MODEL", "").strip()
if env:
return env
try:
conn = psycopg2.connect(**DB_CONFIG)
try:
value = _kv(conn, "app_settings", "SUMMARY_MODEL")
return value or DEFAULT_SUMMARY_MODEL
finally:
conn.close()
except Exception: # noqa: BLE001
return DEFAULT_SUMMARY_MODEL
def resolve_base_url() -> str:
return os.getenv("NOUS_BASE_URL", DEFAULT_NOUS_BASE_URL).strip() or DEFAULT_NOUS_BASE_URL
def call_llm(prompt: str, *, api_key: str, model: str, base_url: str, json_mode: bool = False) -> str:
"""Send a prompt to Nous chat completions and return the text (\"\" on failure)."""
if not api_key:
logger.warning("NOUS_API_KEY not set — skipping LLM call")
return ""
return chat(prompt, api_key=api_key, model=model, base_url=base_url, json_mode=json_mode)
# ── Futures (legacy, gated) ────────────────────────────────────────────────
def fetch_current_futures_prices() -> dict:
"""Live futures prices. Only meaningful when INCLUDE_FUTURES=1."""
if not INCLUDE_FUTURES:
return {}
try:
import yfinance as yf # noqa: PLC0415
except ImportError:
logger.warning(
"INCLUDE_FUTURES=1 but yfinance is not installed — install it to enable futures prices"
)
return {}
prices: dict = {}
for category, tickers in FUTURES_TICKERS.items():
for ticker in tickers:
try:
data = yf.Ticker(ticker).history(period="1d", interval="1m")
if not data.empty:
last_price = data["Close"].iloc[-1]
prices[ticker] = {
"price": round(last_price, 2),
"change_pct": round(
(last_price - data["Open"].iloc[0]) / data["Open"].iloc[0] * 100, 2
) if len(data) > 1 else 0,
"timestamp": datetime.utcnow().strftime("%Y-%m-%d %H:%M UTC"),
"category": category,
}
else:
prices[ticker] = {"price": None, "error": "No data"}
except Exception as exc: # noqa: BLE001
prices[ticker] = {"price": None, "error": str(exc)}
return prices
def build_futures_context() -> str:
ctx = f"CURRENT FUTURES PRICES (as of {datetime.now().strftime('%Y-%m-%d %H:%M UTC')}):\n"
for ticker, info in fetch_current_futures_prices().items():
if info.get("price") is not None:
ctx += (
f"- {ticker} ({info['category']}): ${info['price']:.2f} "
f"({info['change_pct']:+.2f}% today)\n"
)
else:
ctx += f"- {ticker}: unavailable ({info.get('error', 'unknown error')})\n"
return ctx
# ── DB helpers ─────────────────────────────────────────────────────────────
def ensure_tables() -> None:
"""Idempotently create the news tables if missing.
Normally created by alembic 003_news + 005_news_items when the app
container starts, but this summarizer may boot before the app has run
migrations (compose only guarantees `db` is up, not that alembic has
run). Mirrors the scraper pipeline's own CREATE TABLE IF NOT EXISTS so
either start order is safe.
"""
ddl = """
CREATE TABLE IF NOT EXISTS articles (
id SERIAL PRIMARY KEY,
title TEXT,
url TEXT UNIQUE,
content TEXT,
domain TEXT,
timestamp TIMESTAMPTZ
);
CREATE TABLE IF NOT EXISTS article_summaries (
id SERIAL PRIMARY KEY,
summary_text TEXT NOT NULL,
batch_timestamp TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
ALTER TABLE article_summaries ADD COLUMN IF NOT EXISTS model TEXT;
CREATE TABLE IF NOT EXISTS news_items (
id SERIAL PRIMARY KEY,
summary_id INTEGER REFERENCES article_summaries(id) ON DELETE CASCADE,
kind TEXT NOT NULL,
headline TEXT NOT NULL,
importance TEXT NOT NULL,
location_name TEXT,
lat DOUBLE PRECISION,
lon DOUBLE PRECISION,
location_confidence TEXT,
category TEXT,
url TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS ix_news_items_kind_created
ON news_items (kind, created_at DESC);
CREATE INDEX IF NOT EXISTS ix_news_items_map_bbox
ON news_items (lon, lat)
WHERE kind = 'map' AND lat IS NOT NULL AND lon IS NOT NULL;
"""
try:
conn = psycopg2.connect(**DB_CONFIG)
cur = conn.cursor()
cur.execute(ddl)
conn.commit()
cur.close()
conn.close()
except Exception as exc: # noqa: BLE001
logger.error("Error ensuring news tables: %s", exc)
def get_recent_news() -> list[dict]:
"""Fetch articles from the last SUMMARY_WINDOW_HOURS (content > 100 chars)."""
query = """
SELECT title, content, url, domain
FROM articles
WHERE timestamp > NOW() - make_interval(hours => %s)
AND content IS NOT NULL AND length(content) > 100
ORDER BY timestamp DESC;
"""
try:
conn = psycopg2.connect(**DB_CONFIG)
cur = conn.cursor()
cur.execute(query, (SUMMARY_WINDOW_HOURS,))
rows = cur.fetchall()
cur.close()
conn.close()
return [
{"title": r[0], "content": r[1], "url": r[2], "domain": r[3]}
for r in rows
]
except Exception as exc: # noqa: BLE001
logger.error("Database error reading articles: %s", exc)
return []
def _already_summarized_this_hour() -> bool:
"""True when article_summaries already has a row for the current UTC hour."""
if os.getenv("NEWS_SUMMARIZE_FORCE", "") == "1":
return False
query = (
"SELECT 1 FROM article_summaries "
"WHERE batch_timestamp >= date_trunc('hour', NOW() AT TIME ZONE 'utc')"
)
try:
conn = psycopg2.connect(**DB_CONFIG)
cur = conn.cursor()
cur.execute(query)
row = cur.fetchone()
cur.close()
conn.close()
return row is not None
except Exception as exc: # noqa: BLE001
logger.error("Error checking hourly idempotency: %s", exc)
return False
def save_batch(summary_en: str, model: str, ticker: list, map_items: list) -> None:
"""Insert the master brief plus flagged ticker/map rows."""
ticker_rows = select_ticker(ticker or [])
map_rows = select_map(map_items or [])
text = (summary_en or "").strip()
if len(text) < 10 and not ticker_rows and not map_rows:
logger.info("Summary too short or empty. Skipping save.")
return
insert_item = """
INSERT INTO news_items (
summary_id, kind, headline, importance, location_name,
lat, lon, location_confidence, category, url
) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
"""
try:
conn = psycopg2.connect(**DB_CONFIG)
cur = conn.cursor()
cur.execute(
"INSERT INTO article_summaries (summary_text, model) VALUES (%s, %s) RETURNING id",
(text, model),
)
summary_id = cur.fetchone()[0]
for row in ticker_rows:
cur.execute(
insert_item,
(
summary_id,
"ticker",
row.get("headline"),
row.get("importance"),
row.get("location_name"),
None,
None,
None,
None,
row.get("url"),
),
)
for row in map_rows:
cur.execute(
insert_item,
(
summary_id,
"map",
row.get("headline"),
row.get("importance"),
row.get("location_name"),
row.get("lat"),
row.get("lon"),
row.get("location_confidence"),
row.get("category"),
row.get("url"),
),
)
conn.commit()
logger.info(
"Master summary saved id=%s model=%s ticker=%d map=%d",
summary_id, model, len(ticker_rows), len(map_rows),
)
cur.close()
conn.close()
except Exception as exc: # noqa: BLE001
logger.error("Error saving batch to DB: %s", exc)
# ── Orchestration ──────────────────────────────────────────────────────────
def build_map_prompt(batch: list[dict]) -> str:
batch_text = "\n\n".join(
f"Title: {a['title']}\nSource: {a['domain']}\nURL: {a['url']}\nContent: {a['content'][:1500]}"
for a in batch
)
template = map_prompt_template()
prefix = build_futures_context() + "\n" if INCLUDE_FUTURES else ""
try:
return prefix + template.format(batch_text=batch_text)
except KeyError:
return prefix + template
def build_master_prompt(final_input: str) -> str:
template = summary_prompt_template()
prefix = build_futures_context() + "\n" if INCLUDE_FUTURES else ""
try:
return prefix + template.format(final_input=final_input)
except KeyError:
return prefix + template
def summarize_news() -> None:
"""Map-reduce summarize recent articles and store brief + ticker + map."""
ensure_tables()
if _already_summarized_this_hour():
logger.info(
"Skipping summarize: article_summaries already has a row this UTC hour "
"(set NEWS_SUMMARIZE_FORCE=1 to override)"
)
return
api_key = resolve_api_key()
model = resolve_model()
base_url = resolve_base_url()
if not api_key:
logger.warning("NOUS_API_KEY unset in env and api_keys — idle this run")
return
articles = get_recent_news()
if not articles:
logger.info("No new articles found in the last %sh.", SUMMARY_WINDOW_HOURS)
return
logger.info(
"Processing %d articles with %s (batch_size=%d, futures=%s)...",
len(articles), model, BATCH_SIZE, INCLUDE_FUTURES,
)
partial_summaries: list[str] = []
for i in range(0, len(articles), BATCH_SIZE):
batch = articles[i : i + BATCH_SIZE]
logger.info(
"map batch %d/%d (%d articles)",
i // BATCH_SIZE + 1, -(-len(articles) // BATCH_SIZE), len(batch),
)
summary = call_llm(
build_map_prompt(batch),
api_key=api_key,
model=model,
base_url=base_url,
json_mode=False,
)
if summary:
partial_summaries.append(summary)
final_input = "\n\n".join(partial_summaries)
if not final_input.strip():
logger.warning("No partial summaries produced — nothing to reduce.")
return
logger.info("reduce phase over %d partial summaries", len(partial_summaries))
master_raw = call_llm(
build_master_prompt(final_input),
api_key=api_key,
model=model,
base_url=base_url,
json_mode=True,
)
if not master_raw:
logger.warning("Reduce phase returned empty — nothing to persist.")
return
parsed = parse_reduce_json(master_raw)
save_batch(parsed["summary_en"], model, parsed["ticker"], parsed["map_items"])
if __name__ == "__main__":
summarize_news()