all systems operationalIndia

CVE Trove: a vulnerability intelligence pipeline

Pulls 15 public vulnerability feeds on a schedule, merges them into one record per CVE, and streams the result into MongoDB through Celery workers.

state
deploying
stage
growing
status
active
started
2024.02
source
github ↗
topics

The facts about a single vulnerability are scattered across the internet. NVD has the description and CVSS score. CISA says whether it's being exploited right now. EPSS estimates how likely it is to be exploited in the next 30 days. ExploitDB, Metasploit and Nuclei show whether a working exploit already exists. Sigma and Snort have rules to detect it.

To decide what to patch first you need all of that in one place. CVE Trove collects it automatically. Every day it pulls 15 public sources, merges everything it learns about each CVE into one record, and writes the records into MongoDB through a queue of workers. Weekly, it also collects exploit templates, and at startup it loads the MITRE CWE and CAPEC taxonomies.

Source code: github.com/sumit-kumar-03/cve-trove

What it collects

KindSources
Advisories: what the vulnerability isNVD (National Vulnerability Database), the MITRE CVE list, GitHub Advisory Database, OSV, CISA alerts
Exploitation: is it being attacked or attackableCISA Known Exploited Vulnerabilities, ExploitDB, Metasploit, Nuclei templates, PoC-in-GitHub, Tenable PoCs, BadThings
Detection: how to spot an attackSigma rules, Snort community rules
Likelihood: how probable exploitation isEPSS scores
Taxonomy (at startup)MITRE CWE (weakness types) and CAPEC (attack patterns)
Exploit templates (weekly)Nuclei PoC templates, raw Tenable PoC scripts

Some sources are downloads (NVD yearly feeds, the EPSS CSV, the CISA JSON). Others are whole git repositories that it clones and walks (ExploitDB, Metasploit, Nuclei templates, GitHub Advisory Database).

How it works

cron (in the worker container)
 ├─ daily 00:00   run_cve_update.py ─▶ CveTrove
 │                    for each registered scraper:  fetch_data() → parse_data()
 │                        └─ merges into ResultStore  { "CVE-2024-3094": Cve(...), ... }
 │                    flush() ─▶ batches of 1,000 ─▶ RabbitMQ queue "cve_trove_task"
 │                                                      │
 ├─ weekly Sun 03:00  run_poc_update.py ─▶ PocTrove ────┤
 └─ at startup        run_atlas_update.py (CWE/CAPEC) ──┤
                                                        ▼
                                   Celery workers (4 concurrent, late acks)
                                                        │  Beanie BulkWriter
                                                        ▼
                                   MongoDB  cve_trove.CVE  (one document per CVE)

Each CVE becomes one document. Every source adds its piece to the same record:

class Cve(BaseModel):
    cve: str                          # "CVE-2024-3094"
    sources: List[str]                # which feeds mentioned it
    advisories: AdvisoryData          # per-source title, description, CVSS, CPEs, references
    exploits: Exploits                # exploit_available + links to exploit code
    zeroday: Zeroday                  # flagged by CISA alerts / KEV, with references
    epss_score: Optional[float]       # probability of exploitation (EPSS)
    priority_score: Optional[float]
    taxonomy: Taxonomy                # CWE / CAPEC mappings
    remediations: List[Remediation]

The design choices that matter

  • Scrapers are plugins. Each source is one class that inherits BaseScraper and registers itself with a decorator. Adding a source means adding one file; nothing else changes.
  • Merge in memory, write in batches. Scrapers write into a shared in-memory dictionary keyed by CVE ID, so 15 sources produce one merged record per CVE rather than 15 fragments. Only then is it flushed, 1,000 records at a time.
  • A queue between scraping and storage. Batches go to RabbitMQ and Celery workers write them to MongoDB. Slow database writes never block scraping, and writes can scale by adding workers.
  • Workers that don't lose work. Tasks are acknowledged only after they finish (task_acks_late) and are requeued if a worker dies mid-task (task_reject_on_worker_lost). Each worker process is recycled after 1,000 tasks to keep memory in check.
  • Shared infrastructure. MongoDB and RabbitMQ live in a separate infra-hub Docker Compose stack on a shared network, so several projects can use one database and one broker.

Run it

It needs Docker and several gigabytes of free disk, because sources such as Metasploit, ExploitDB and the GitHub Advisory Database are cloned as full git repositories.

1. Start the shared infrastructure. MongoDB, RabbitMQ and mongo-express run from the infra-hub repo, which also creates the infra-net network:

git clone https://github.com/sumit-kumar-03/infra-hub.git
cd infra-hub && cp .env.example .env    # set real passwords
docker compose up -d

2. Configure and start the worker. The repo ships a dev.env.example listing every setting. Copy it and use the same credentials you set in infra-hub:

git clone https://github.com/sumit-kumar-03/cve-trove.git
cd cve-trove
cp dev.env.example dev.env      # RabbitMQ + MongoDB credentials, optional GitHub token
docker compose up -d --build
docker compose logs -f cve-trove-worker

On start, the worker loads the CWE and CAPEC taxonomy, installs the cron schedule, and starts four Celery worker processes. The first full CVE run happens at the next midnight in the container's clock, which is usually UTC. To run one immediately:

docker compose exec cve-trove-worker python3 scripts/run_cve_update.py

The repo's GUIDE.md covers each setting, the schedule, more queries, and how to add a scraper.

3. Query the results. Browse them in mongo-express on port 8081, or query directly. For example, CVEs with a public exploit and a high chance of exploitation:

use cve_trove
db.CVE.find(
  { "data.exploits.exploit_available": true, "data.epss_score": { $gt: 0.5 } },
  { _id: 1, "data.epss_score": 1, "data.sources": 1 }
).sort({ "data.epss_score": -1 }).limit(20)

Build it yourself, step by step

This section rebuilds the core ideas in plain Python, so you can adapt them to any "collect from many sources and merge" problem.

1. Define the record every source feeds

Pydantic gives you validation and easy conversion to JSON. Start small and grow the model as you add sources:

from typing import Dict, List, Optional
from pydantic import BaseModel, Field, RootModel

class Exploits(BaseModel):
    exploit_available: bool = False
    references: List[str] = Field(default_factory=list)

class Cve(BaseModel):
    cve: str
    sources: List[str] = Field(default_factory=list)
    epss_score: Optional[float] = None
    exploits: Exploits = Field(default_factory=Exploits)

class CveObject(RootModel):
    root: Dict[str, Cve]          # { "CVE-2024-3094": Cve(...) }

2. A base class that fixes the shape of every scraper

Every source does the same two things: get raw data, then turn it into records. The base class fixes that order and collects errors instead of crashing the whole run:

class BaseScraper:
    source: str = "unknown"

    def __init__(self):
        self.errors: List[str] = []
        self.data: Optional[CveObject] = None

    def fetch_data(self): raise NotImplementedError
    def parse_data(self): raise NotImplementedError

    def run(self, store) -> bool:
        try:
            self.data = store.cves          # every scraper writes into the same dict
            self.fetch_data()
            self.parse_data()
        except Exception as e:
            self.errors.append(str(e))      # one broken feed shouldn't stop the others
        return not self.errors

3. Register scrapers with a decorator

A registry means the orchestrator never needs a hard-coded list. Importing a scraper's module is enough to add it:

from threading import Lock

class ScraperRegistry:
    _scrapers: dict = {}
    _lock = Lock()

    @classmethod
    def register(cls, name):
        def decorator(scraper_cls):
            with cls._lock:
                cls._scrapers.setdefault(name, scraper_cls)
            return scraper_cls
        return decorator

registry = ScraperRegistry()

4. Write a scraper

EPSS publishes one gzipped CSV of scores for every CVE, which makes it a good first source:

import csv, gzip, io, requests

@registry.register("epss")
class Epss(BaseScraper):
    source = "epss"
    URL = "https://epss.cyentia.com/epss_scores-current.csv.gz"

    def fetch_data(self):
        raw = gzip.decompress(requests.get(self.URL, timeout=120).content).decode()
        lines = raw.splitlines()[1:]                 # first line is a comment
        self.rows = list(csv.DictReader(io.StringIO("\n".join(lines))))

    def parse_data(self):
        for row in self.rows:
            cve_id = row["cve"]
            record = self.data.root.setdefault(cve_id, Cve(cve=cve_id))
            record.epss_score = float(row["epss"])
            record.sources.append(self.source)

setdefault is the whole merge strategy. Whichever scraper sees a CVE first creates its record, and every later scraper adds to that same record.

5. Merge in memory, flush in batches

class ResultStore:
    BATCH = 1000

    def __init__(self):
        self.cves = CveObject(root={})

    def flush(self):
        records = [c.model_dump() for c in self.cves.root.values()]
        for i in range(0, len(records), self.BATCH):
            save_cves.delay(records[i : i + self.BATCH])     # a Celery task (next step)

6. Queue the writes with Celery and RabbitMQ

The scraper process only enqueues batches, and workers do the database writes. Late acknowledgement makes the queue at-least-once, so a crashed worker's batch is retried rather than lost. Writing by CVE ID makes retries safe, because the same batch twice leaves the same result:

from celery import Celery
from pymongo import MongoClient, ReplaceOne

app = Celery("cve_trove", broker="amqp://user:pass@infra-rabbitmq:5672//")
app.conf.task_acks_late = True
app.conf.task_reject_on_worker_lost = True
app.conf.worker_max_tasks_per_child = 1000

mongo = MongoClient("mongodb://user:pass@infra-mongodb:27017")

@app.task(queue="cve_trove_task")
def save_cves(batch):
    mongo.cve_trove.CVE.bulk_write(
        [ReplaceOne({"_id": r["cve"]}, {"_id": r["cve"], "data": r}, upsert=True) for r in batch]
    )

CVE Trove itself uses Beanie, an async ODM on top of Motor. Plain pymongo keeps this example short.

7. Orchestrate

import my_scrapers   # importing the package runs every @registry.register

def run():
    store = ResultStore()
    for name, scraper_cls in registry._scrapers.items():
        ok = scraper_cls().run(store)
        print(f"{name}: {'ok' if ok else 'failed'}")
    store.flush()

8. Schedule it inside the container

One container runs both the schedule and the workers. The start script installs the crontab, starts cron in the background, and then runs Celery in the foreground so the container stays up:

# scripts/crontab
0 0 * * *  cd /usr/src/app && python3 scripts/run_cve_update.py >> /var/log/cve_update.log 2>&1
0 3 * * 0  cd /usr/src/app && python3 scripts/run_poc_update.py >> /var/log/poc_update.log 2>&1
#!/bin/sh
set -e
python3 scripts/run_atlas_update.py          # load CWE / CAPEC once at startup
crontab scripts/crontab && crond -b -l 2     # schedule in the background
exec celery --app=cve_trove.celery_service.app.cve_trove_celery_app worker \
     --queues=cve_trove_task --concurrency=4 -Ofair --max-tasks-per-child=1000

9. Share infrastructure across projects

Instead of every project bundling its own database and broker, one stack owns them and creates a network. Other projects join that network as external and reach services by name:

# cve-trove/docker-compose.yml
services:
  cve-trove-worker:
    build: .
    env_file: [dev.env]
    entrypoint: ["sh", "/usr/src/app/scripts/start_worker_with_cron.sh"]
    networks: [infra-net]
    restart: unless-stopped

networks:
  infra-net:
    name: infra-net
    external: true        # created by the infra-hub stack

Ideas to extend it

  • A priority score. Combine CVSS, EPSS and "listed by CISA as exploited" into one number and store it in the existing priority_score field.
  • Only process what changed. Keep the last commit hash of each cloned source and only parse files changed since then, instead of re-walking whole repositories.
  • An API. A small FastAPI service in front of the CVE collection, for example GET /cve/{id} or GET /cves?exploited=true&min_epss=0.5.
  • Per-source health. Record each scraper's success, duration and error list per run, and alert when a feed starts failing.

Connected

shares a topic with this project