Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions common/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,7 @@ def __init__(self):
self.insights_load_rule_cache_ttl_sec = int(
os.getenv("INSIGHTS_LOAD_RULE_CACHE_TTL_SEC", "0")
) # 0 -> local cache disabled -> always use DB cache
self.insights_cve_cache_ttl_sec = int(os.getenv("INSIGHTS_CVE_CACHE_TTL_SEC", "0"))


# pylint: disable=too-many-instance-attributes
Expand Down
1 change: 1 addition & 0 deletions common/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

APP_VERSION = "2.79.0"
TIMESTAMP_LAST_REPO_BASED_EVAL = "last_eval_repo_based"
TIMESTAMP_LAST_CVE_SYNC = "last_cve_sync"
VMAAS_CVES_ENDPOINT = f"{CFG.vmaas_host}/api/vmaas/v3/cves"
VMAAS_REPOS_ENDPOINT = f"{CFG.vmaas_host}/api/vmaas/v3/repos"
VMAAS_OS_ENDPOINT = f"{CFG.vmaas_host}/api/vmaas/v3/os/vulnerability/report"
Expand Down
1 change: 1 addition & 0 deletions conf/evaluator.env
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,4 @@ MAX_LOADED_EVALUATOR_MSGS=20
USE_VMAAS_GO=true
INVENTORY_VIEWS_TOPIC=platform.inventory.host-apps
INSIGHTS_LOAD_RULE_CACHE_TTL_SEC=360
INSIGHTS_CVE_CACHE_TTL_SEC=360
7 changes: 7 additions & 0 deletions deploy/clowdapp.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -417,6 +417,8 @@ objects:
value: ${DB_STATEMENT_TIMEOUT}
- name: INSIGHTS_LOAD_RULE_CACHE_TTL_SEC
value: ${INSIGHTS_LOAD_RULE_CACHE_TTL_SEC}
- name: INSIGHTS_CVE_CACHE_TTL_SEC
value: ${INSIGHTS_CVE_CACHE_TTL_SEC}
resources:
limits:
cpu: ${CPU_LIMIT_EVALUATOR_RECALC}
Expand Down Expand Up @@ -467,6 +469,8 @@ objects:
value: ${DB_STATEMENT_TIMEOUT}
- name: INSIGHTS_LOAD_RULE_CACHE_TTL_SEC
value: ${INSIGHTS_LOAD_RULE_CACHE_TTL_SEC}
- name: INSIGHTS_CVE_CACHE_TTL_SEC
value: ${INSIGHTS_CVE_CACHE_TTL_SEC}
resources:
limits:
cpu: ${CPU_LIMIT_EVALUATOR_UPLOAD}
Expand Down Expand Up @@ -977,6 +981,9 @@ parameters:
- name: INSIGHTS_LOAD_RULE_CACHE_TTL_SEC
description: TTL in seconds for evaluator insights_rule cache (0 disables caching)
value: "360"
- name: INSIGHTS_CVE_CACHE_TTL_SEC
description: TTL in seconds for evaluator cve_cache cache refresh (0 to refresh after each vmaas_sync)
value: "7200" # 2 hours
- name: MAX_LOADED_LISTENER_MSGS
value: "200"
- name: FLOORIST_SUSPEND
Expand Down
53 changes: 48 additions & 5 deletions evaluator/logic.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
from psycopg.types.json import Jsonb
from psycopg_pool.pool_async import AsyncConnectionPool

from common.constants import TIMESTAMP_LAST_CVE_SYNC
from common.constants import format_vmaas_cve_endpoint
from common.logging import get_logger
from common.peewee_model import VulnerabilityState
Expand Down Expand Up @@ -58,7 +59,10 @@ def __init__(self, db_pool: AsyncConnectionPool):
self.vulnerable_package_cache: Dict[(int, int, Optional[int]), VulnerablePackageCache] = {}
self.skipped_rules = ["CVE_2017_5715_cpu_virt|VIRT_CVE_2017_5715_CPU_3_ONLYKERNEL", "CVE_2017_5715_cpu_virt"]
self.rule_cache: Dict[str, RuleCache] = {}
self.rule_cache_expires_at: Optional[datetime] = None
self.cache_metadata = {
"rule": {"expires_at": None},
"cve": {"loaded_at": None, "expires_at": None},
}

async def init(self):
"""Async constructor"""
Expand All @@ -80,6 +84,43 @@ async def _load_cve_impact_cache(self) -> Dict[str, CveImpactCache]:
cache[cve_impact["name"]] = CveImpactCache(cve_impact["id"])
return cache

def _cve_cache_expired(self) -> bool:
"""Check if CVE cache TTL expired"""
if CFG.insights_cve_cache_ttl_sec <= 0:
return True # TTL disabled, check DB on every evaluation (if vmaas_sync ran)
if self.cache_metadata["cve"]["expires_at"] is None:
return True
return datetime.now(timezone.utc) >= self.cache_metadata["cve"]["expires_at"]

async def _get_last_cve_sync_tms(self) -> Optional[datetime]:
"""Query timestamp_kv for last CVE sync timestamp"""
async with self.db_pool.connection() as conn:
async with conn.cursor(row_factory=dict_row) as cur:
await cur.execute("SELECT value FROM timestamp_kv WHERE name = %s", (TIMESTAMP_LAST_CVE_SYNC,))
row = await cur.fetchone()
return row["value"] if row else None

def _set_cve_cache_expiry(self) -> None:
ttl = CFG.insights_cve_cache_ttl_sec
if ttl > 0:
self.cache_metadata["cve"]["expires_at"] = datetime.now(timezone.utc) + timedelta(seconds=ttl)
else:
self.cache_metadata["cve"]["expires_at"] = None

async def _refresh_cve_cache(self) -> None:
"""Refresh CVE cache if vmaas_sync ran since last load and TTL expired"""
if self._cve_cache_expired():
last_sync = await self._get_last_cve_sync_tms()
# Reload only if vmaas_sync ran since we last loaded
if last_sync and (not self.cache_metadata["cve"]["loaded_at"] or last_sync > self.cache_metadata["cve"]["loaded_at"]):
LOGGER.info("CVE cache stale (vmaas_sync ran at %s), reloading", last_sync)
self.cve_cache = await self._load_cve_cache()
self.cache_metadata["cve"]["loaded_at"] = datetime.now(timezone.utc)
# We can possibly reset TTL only after cache reload (at this place) to make the cache more resilient to changes

# Reset TTL regardless to avoid DB query on next N evaluations (polling)
self._set_cve_cache_expiry()

async def _load_cve_cache(self) -> Dict[str, CveCache]:
"""Load cve cache from DB"""
cache = {}
Expand All @@ -103,17 +144,17 @@ def _set_rule_cache_expiry(self) -> None:
"""Set rule cache expiry; only used when TTL > 0 (see _rule_cache_expired)"""
ttl = CFG.insights_load_rule_cache_ttl_sec
if ttl > 0:
self.rule_cache_expires_at = datetime.now(timezone.utc) + timedelta(seconds=ttl)
self.cache_metadata["rule"]["expires_at"] = datetime.now(timezone.utc) + timedelta(seconds=ttl)
else:
self.rule_cache_expires_at = None
self.cache_metadata["rule"]["expires_at"] = None

def _rule_cache_expired(self) -> bool:
"""TTL <= 0 disables caching (always reload on get from DB)"""
if CFG.insights_load_rule_cache_ttl_sec <= 0:
return True
if not self.rule_cache or self.rule_cache_expires_at is None:
if not self.rule_cache or self.cache_metadata["rule"]["expires_at"] is None:
return True
return datetime.now(timezone.utc) >= self.rule_cache_expires_at
return datetime.now(timezone.utc) >= self.cache_metadata["rule"]["expires_at"]

async def _refresh_rule_cache(self, conn: Optional[AsyncConnection] = None) -> None:
"""Refresh self.rule_cache from DB when TTL expired (see _load_rule_cache)"""
Expand Down Expand Up @@ -516,6 +557,8 @@ async def _evaluate_vmaas_res(
) -> dict:
"""Insert vmaas cve results"""

await self._refresh_cve_cache()

with EVAL_PART_TIME.labels(part="get_or_upsert").time():
# system is potentially vulnerable to cves returned from vmaas
for cve_adv in playbook_cves:
Expand Down
198 changes: 198 additions & 0 deletions tests/common_tests/test_evaluator_cve_cache.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,198 @@
# pylint: disable=missing-docstring,redefined-outer-name
"""Evaluator CVE cache TTL-based refresh"""

from contextlib import asynccontextmanager
from datetime import datetime
from datetime import timedelta
from datetime import timezone

import pytest

pytestmark = pytest.mark.asyncio(loop_scope="function")


class FakeAsyncPool:

@asynccontextmanager
async def connection(self):
yield object()


async def test_cve_cache_expired_returns_true_when_no_expiry_set():
"""Cache with no expiry timestamp is considered expired."""
import evaluator.logic as logic_module
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())
logic.cache_metadata["cve"]["expires_at"] = None

# Mock TTL as enabled
import unittest.mock

with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 300):
assert logic._cve_cache_expired() is True


async def test_cve_cache_expired_returns_false_when_ttl_disabled():
"""TTL <= 0 autorefreshes regardless ttl"""
import unittest.mock

import evaluator.logic as logic_module
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())

with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 0):
assert logic._cve_cache_expired() is True


async def test_cve_cache_expired_returns_true_when_ttl_past():
"""Cache is expired when current time > expiry time."""
import unittest.mock

import evaluator.logic as logic_module
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())
logic.cache_metadata["cve"]["expires_at"] = datetime.now(timezone.utc) - timedelta(seconds=1)

with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 300):
assert logic._cve_cache_expired() is True


async def test_refresh_cve_cache_skips_when_ttl_not_expired():
"""When TTL hasn't expired, skip DB query and reload."""
import unittest.mock

import evaluator.logic as logic_module
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())
future_expiry = datetime.now(timezone.utc) + timedelta(seconds=300)
logic.cache_metadata["cve"]["expires_at"] = future_expiry

# Mock TTL as enabled (> 0)
with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 300):
# Should return early without any DB call
await logic._refresh_cve_cache()

# Expiry should remain unchanged since TTL not expired
assert logic.cache_metadata["cve"]["expires_at"] == future_expiry


async def test_refresh_cve_cache_reloads_when_sync_ran():
"""When vmaas_sync ran after last cache load, reload cache."""
import unittest.mock

import evaluator.logic as logic_module
from evaluator.common import CveCache
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())

# Cache was loaded at T0
cache_loaded_at = datetime(2026, 1, 1, 10, 0, 0, tzinfo=timezone.utc)
logic.cache_metadata["cve"]["loaded_at"] = cache_loaded_at

# vmaas_sync ran at T1 (after T0)
vmaas_sync_time = datetime(2026, 1, 1, 10, 30, 0, tzinfo=timezone.utc)

# Track if cache was reloaded
reload_called = []

async def fake_get_last_sync():
return vmaas_sync_time

async def fake_load_cache():
reload_called.append(1)
return {"CVE-2026-NEW": CveCache(999, 1, False)}

logic._get_last_cve_sync_tms = unittest.mock.AsyncMock(side_effect=fake_get_last_sync)
logic._load_cve_cache = unittest.mock.AsyncMock(side_effect=fake_load_cache)

with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 300):
await logic._refresh_cve_cache()

# Cache should have been reloaded
assert len(reload_called) == 1
# loaded_at should be updated to now (approximately)
assert logic.cache_metadata["cve"]["loaded_at"] > cache_loaded_at
# TTL should be reset (initially its None)
assert logic.cache_metadata["cve"]["expires_at"] is not None


async def test_refresh_cve_cache_skips_reload_when_sync_not_ran():
"""When vmaas_sync hasn't run since last load, skip reload but reset TTL."""
import unittest.mock

import evaluator.logic as logic_module
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())

# Cache was loaded at T1
cache_loaded_at = datetime(2026, 1, 1, 11, 0, 0, tzinfo=timezone.utc)
logic.cache_metadata["cve"]["loaded_at"] = cache_loaded_at

# TTL expired 5 minutes ago
old_expires_at = datetime.now(timezone.utc) - timedelta(minutes=5)
logic.cache_metadata["cve"]["expires_at"] = old_expires_at

# vmaas_sync last ran at T0 (before T1)
vmaas_sync_time = datetime(2026, 1, 1, 10, 0, 0, tzinfo=timezone.utc)

# Track if cache was reloaded
reload_called = []

async def fake_get_last_sync():
return vmaas_sync_time

async def fake_load_cache():
reload_called.append(1)
return {}

logic._get_last_cve_sync_tms = unittest.mock.AsyncMock(side_effect=fake_get_last_sync)
logic._load_cve_cache = unittest.mock.AsyncMock(side_effect=fake_load_cache)

with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 300):
await logic._refresh_cve_cache()

# Cache should NOT have been reloaded
assert len(reload_called) == 0
# loaded_at should remain unchanged
assert logic.cache_metadata["cve"]["loaded_at"] == cache_loaded_at
# TTL should be reset (polling mode - check again after TTL period)
new_expires_at = logic.cache_metadata["cve"]["expires_at"]
assert new_expires_at is not None
assert new_expires_at != old_expires_at # Verify it actually changed
assert new_expires_at > datetime.now(timezone.utc) # Should be in the future


async def test_refresh_cve_cache_resets_ttl_when_initially_none():
"""TTL should be set even when initially None (first check after startup)."""
import unittest.mock

import evaluator.logic as logic_module
from evaluator.logic import EvaluatorLogic

logic = EvaluatorLogic(FakeAsyncPool())

# Simulate startup state: cache loaded but TTL never set
cache_loaded_at = datetime(2026, 1, 1, 11, 0, 0, tzinfo=timezone.utc)
logic.cache_metadata["cve"]["loaded_at"] = cache_loaded_at
logic.cache_metadata["cve"]["expires_at"] = None # Never set before

# vmaas_sync last ran before cache load (no reload needed)
vmaas_sync_time = datetime(2026, 1, 1, 10, 0, 0, tzinfo=timezone.utc)

async def fake_get_last_sync():
return vmaas_sync_time

logic._get_last_cve_sync_tms = unittest.mock.AsyncMock(side_effect=fake_get_last_sync)

with unittest.mock.patch.object(logic_module.CFG, "insights_cve_cache_ttl_sec", 300):
await logic._refresh_cve_cache()

# TTL should now be set
assert logic.cache_metadata["cve"]["expires_at"] is not None
assert logic.cache_metadata["cve"]["expires_at"] > datetime.now(timezone.utc)
6 changes: 3 additions & 3 deletions tests/common_tests/test_evaluator_rule_cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ async def fake_load_rule_cache(_self, _conn):


async def test_refresh_rule_cache_reloads_after_ttl_expires(monkeypatch):
"""TTL > 0: after rule_cache_expires_at is in the past, the next refresh runs a full load again"""
"""TTL > 0: after logic.cache_metadata["rule"]["expires_at"] is in the past, the next refresh runs a full load again"""
import evaluator.logic as logic_module
from evaluator.common import RuleCache
from evaluator.logic import EvaluatorLogic
Expand All @@ -93,7 +93,7 @@ async def fake_load_rule_cache(_self, _conn):
await logic._refresh_rule_cache()
assert len(loads) == 1

logic.rule_cache_expires_at = datetime.now(timezone.utc) - timedelta(seconds=1)
logic.cache_metadata["rule"]["expires_at"] = datetime.now(timezone.utc) - timedelta(seconds=1)

await logic._refresh_rule_cache()
assert len(loads) == 2
Expand Down Expand Up @@ -155,7 +155,7 @@ async def full_load_other_rules_only(_self, _conn):

assert len(full_load_calls) == 1
assert rule_id not in logic.rule_cache
assert logic.rule_cache_expires_at is not None
assert logic.cache_metadata["rule"]["expires_at"] is not None

# Empty sys_vuln_rows: rule-only hit branch (no prior VMAAS row for this CVE)
out = await logic._evaluate_advisor_res(rule_results, {}, platform, set(), conn)
Expand Down
15 changes: 15 additions & 0 deletions vmaas_sync/vmaas_sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,15 @@ def _set_last_repobased_eval_tms(cur, timestamp: dt.datetime):
return ret


def _set_last_cve_sync_tms(cur, timestamp: dt.datetime):
"""Update last CVE sync timestamp"""
cur.execute(
"""insert into timestamp_kv (name, value) values (%s, %s)
on conflict (name) do update set value = %s""",
(constants.TIMESTAMP_LAST_CVE_SYNC, timestamp, timestamp),
)


def _get_updated_repos(conn) -> dict[str, list[str]]:
"""
Get repos updated since last repo-based evaluation
Expand Down Expand Up @@ -341,6 +350,12 @@ def sync_cve_md():
)

conn.commit()

# Record sync completion timestamp for evaluator cache refresh
with conn.cursor() as cur:
_set_last_cve_sync_tms(cur, dt.datetime.now(dt.timezone.utc))
conn.commit()

LOGGER.info("Finished syncing CVE metadata")
return True

Expand Down
Loading