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.
432 lines
16 KiB
Python
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()
|