diff --git a/common/config.py b/common/config.py index f5d194c17..1f2250f87 100644 --- a/common/config.py +++ b/common/config.py @@ -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 diff --git a/common/constants.py b/common/constants.py index 92e53c965..899761e8b 100644 --- a/common/constants.py +++ b/common/constants.py @@ -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" diff --git a/conf/evaluator.env b/conf/evaluator.env index fc116b3c3..980f7d60f 100644 --- a/conf/evaluator.env +++ b/conf/evaluator.env @@ -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 diff --git a/deploy/clowdapp.yaml b/deploy/clowdapp.yaml index 28f18d638..79611a7a2 100644 --- a/deploy/clowdapp.yaml +++ b/deploy/clowdapp.yaml @@ -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} @@ -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} @@ -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 diff --git a/evaluator/logic.py b/evaluator/logic.py index 0236b1e78..c0665c5d3 100644 --- a/evaluator/logic.py +++ b/evaluator/logic.py @@ -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 @@ -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""" @@ -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 = {} @@ -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)""" @@ -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: diff --git a/tests/common_tests/test_evaluator_cve_cache.py b/tests/common_tests/test_evaluator_cve_cache.py new file mode 100644 index 000000000..438ba8384 --- /dev/null +++ b/tests/common_tests/test_evaluator_cve_cache.py @@ -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) diff --git a/tests/common_tests/test_evaluator_rule_cache.py b/tests/common_tests/test_evaluator_rule_cache.py index 8375345f1..f62eac652 100644 --- a/tests/common_tests/test_evaluator_rule_cache.py +++ b/tests/common_tests/test_evaluator_rule_cache.py @@ -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 @@ -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 @@ -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) diff --git a/vmaas_sync/vmaas_sync.py b/vmaas_sync/vmaas_sync.py index 026631f98..0311a9f30 100644 --- a/vmaas_sync/vmaas_sync.py +++ b/vmaas_sync/vmaas_sync.py @@ -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 @@ -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