­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ ­ """ This program is free software: you can redistribute it and/or modify it under the terms of the GNU General Public License as published by the Free Software Foundation, either version 3 of the License, or (at your option) any later version. This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the GNU General Public License for more details. You should have received a copy of the GNU General Public License along with this program.  If not, see . Copyright © 2019 Cloud Linux Software Inc. This software is also available under ImunifyAV commercial license, see """ import grp import logging import os import pwd import re from collections import defaultdict from pathlib import Path from typing import Any, Dict, Iterable, List, Literal, Optional, Set, Tuple from datetime import datetime, timedelta, timezone from peewee import Case, fn from defence360agent.subsys.panels import hosting_panel from defence360agent.utils import get_results_iterable_expression, is_cluster from defence360agent.utils.threads import to_thread from defence360agent.model.analyst_cleanup import AnalystCleanupRequest from defence360agent.api.server.analyst_cleanup import AnalystCleanupAPI from imav.malwarelib.config import ( MalwareHitStatus, MalwareScanResourceType, QueuedScanState, ) from defence360agent.rpc_tools.validate import ValidationError from imav.malwarelib.model import MalwareHit, MalwareScan from imav.malwarelib.scan.crontab import is_crontab from imav.malwarelib.tenant_path import ( TenantPath, split_marked, split_prefixed, to_prefixed, ) from imav.malwarelib.utils.cloudways import CloudwaysUser logger = logging.getLogger(__name__) async def registered_apps() -> set: """Snapshot of registered k8s application ids (panel users).""" return {u["user"] for u in await panel_users()} async def validate_k8s_user( user: Optional[str], registered: Optional[set] = None ) -> None: """Reject a missing or unregistered K8s app ID. An unknown user (typo, decommissioned app) would otherwise enter the scan/submit machinery and hit an app_id with no subscriber. Reject at the input boundary, before any queue entry or scan dir exists. ``registered``: optional pre-fetched snapshot (avoids re-enumerating panel users per item in bulk loops). """ if not user: raise ValidationError("--user should be specified for k8s apps") if registered is None: registered = await registered_apps() if user not in registered: raise ValidationError( "--user %r is not a registered K8s application" % user ) async def resolve_tenant_path( path: str, user: Optional[str] = None, *, require_tenant: bool = False, registered: Optional[set] = None, ) -> TenantPath: """Input-boundary parser accepting both supported path forms on k8s: a bare path plus an explicit ``--user``, or a tenant-prefixed path ``//var/www/...``. - explicit ``user`` wins: it is validated against registered apps and the path is treated as bare (a redundant ``//`` prefix is stripped, so passing both forms together stays idempotent); - without ``user``, the first path component is taken as the tenant iff it is a registered application, else the path is bare; - bare path with no tenant raises ValidationError when ``require_tenant`` is set; - outside k8s the path is returned untouched (never treated as prefixed). """ if not is_cluster(): return TenantPath(path, user=user) if registered is None: registered = await registered_apps() if user is not None: await validate_k8s_user(user, registered) candidate, _ = split_prefixed(path) if candidate and candidate != user and candidate in registered: logger.warning( "path %r looks prefixed with app %r but --user %r was " "given explicitly; treating the path as bare", path, candidate, user, ) return TenantPath.from_prefixed(path, user=user) candidate, bare = split_prefixed(path) if candidate and candidate in registered: return TenantPath(bare, user=candidate) if require_tenant: raise ValidationError( "On Kubernetes the path must carry the application id as its" " first component (//var/www/...) or --user must be" " given; %r does not resolve to a registered application" % path ) return TenantPath(path) async def split_stored_paths(items): """Present stored entries as humans read them: the in-container path, with the tenant it is scoped to in its own field. The registry decides what counts as a tenant — a host path must not lose its first directory to a syntactic guess. """ applications = await registered_apps() for item in items: if not item.get("path"): continue tenant, bare = split_marked(item["path"], applications) if tenant is not None: item["user"], item["path"] = tenant, bare return items def stub_entry(): return { "user": None, "home": None, "infected": 0, "infected_db": 0, "_infected_total": 0, "scan_id": None, "scan_date": None, "scan_status": None, "cleanup_status": None, "analyst_status": None, } def system_users(): """ Get all system users and initialize a dict for them. If a user has leftover config files after being deleted then the panel API might treat him as existent. This is resolved by checking that a system user is a panel user. """ for entry in pwd.getpwall(): u = stub_entry() u["user"] = entry.pw_name u["home"] = entry.pw_dir yield u async def panel_users(): """ Get panel users with their home directories. Home directories are determined by: 1. Integration script (if panel provides get_user_homes()) 2. Fallback to /etc/passwd Users from integration script that provide home but aren't in /etc/passwd are also included, enabling decoupled environments. """ panel = hosting_panel.HostingPanel() panel_user_names = set(await panel.get_users()) # Get home directories from integration script (if available) try: user_homes = await panel.get_user_homes() except Exception: user_homes = {} # Build from /etc/passwd, override home if integration provides it result = [] passwd_users = set() for sys_user in system_users(): if sys_user["user"] in panel_user_names: passwd_users.add(sys_user["user"]) # Prefer home from integration script if available if sys_user["user"] in user_homes and user_homes[sys_user["user"]]: sys_user["home"] = user_homes[sys_user["user"]] result.append(sys_user) # Add users from integration script that aren't in passwd (if they have home) for username, home in user_homes.items(): if ( username in panel_user_names and username not in passwd_users and home ): u = stub_entry() u["user"] = username u["home"] = home result.append(u) return result def get(user_list, **kwargs) -> dict: for u in user_list: if all([u[k] == v for k, v in kwargs.items()]): return u return stub_entry() def _latest_completed_scan_per_app(user_list): """Latest completed scan per app (K8s): the tenant-prefixed scan path encodes the app id, so home paths are unique per app even when apps share the same in-container homedir. Yields (user, scan) pairs.""" paths = { to_prefixed(u["home"], u["user"]): u["user"] for u in user_list if u.get("home") and u.get("user") } def expr(_paths): return ( MalwareScan.select( MalwareScan.scanid, MalwareScan.completed, MalwareScan.path ) .where(MalwareScan.path.in_(_paths)) .group_by(MalwareScan.path) .having(MalwareScan.completed == fn.Max(MalwareScan.completed)) ) for scan in get_results_iterable_expression(expr, list(paths)): yield paths[scan.path], scan def update_infected_count_and_last_scan(user_list): homes = [u["home"] for u in user_list] def expr(_homes): q = ( MalwareScan.select( MalwareScan.scanid, MalwareScan.completed, MalwareScan.path ) .group_by(MalwareScan.path) .having(MalwareScan.completed == fn.Max(MalwareScan.completed)) .where(MalwareScan.path.in_(_homes)) ) return q # FIXME: refactor this (lots of duplication) grouped_hits = ( MalwareHit.select(MalwareHit.user, fn.COUNT().alias("infected")) .where( MalwareHit.is_infected() & (MalwareHit.resource_type == MalwareScanResourceType.FILE.value) ) .group_by(MalwareHit.user) ) grouped_db_hits = ( MalwareHit.select(MalwareHit.user, fn.COUNT().alias("infected_db")) .where( MalwareHit.is_infected() & (MalwareHit.resource_type == MalwareScanResourceType.DB.value) ) .group_by(MalwareHit.user) ) grouped_hits_dict = {entry.user: entry.infected for entry in grouped_hits} grouped_db_hits_dict = { entry.user: entry.infected_db for entry in grouped_db_hits } if is_cluster(): # Apps share a homedir, so a bare path can't pick out one user. # Counts are keyed by the file owner already, and the latest # completed scan is attributed per app via its prefixed path. for u in user_list: u["infected"] = grouped_hits_dict.get(u["user"], 0) for user, entry in _latest_completed_scan_per_app(user_list): u = get(user_list, user=user) u["scan_status"] = QueuedScanState.stopped.value u["scan_id"] = entry.scanid else: for entry in get_results_iterable_expression(expr, homes): u = get(user_list, home=entry.path) u["infected"] = grouped_hits_dict.get(u["user"], 0) u["scan_status"] = QueuedScanState.stopped.value u["scan_id"] = entry.scanid for user, infected_db in grouped_db_hits_dict.items(): u = get(user_list, user=user) u["infected_db"] = infected_db for u in user_list: u["_infected_total"] = u["infected"] + u["infected_db"] def update_running_scan_status(user_list, get_scans): paths = [u["home"] for u in user_list] k8s = is_cluster() for scan, status in get_scans(paths): # `start-user` scans carry the app in args["initiator"], not in # `user`; honour both so shared-home queued scans hit the right app. scan_user = getattr(scan, "user", None) or ( getattr(scan, "args", None) or {} ).get("initiator") if k8s and scan_user: u = get(user_list, user=scan_user) else: u = get(user_list, home=scan.path) u["scan_id"] = scan.scanid u["scan_status"] = status u["scan_type"] = scan.scan_type def update_cleanup_status(user_list): """ Updates cleanup status for the list of panel users If at least on cleanup is running for user then status is 'running' Else if there are any finished cleanups then status is 'stopped' If no started and finished cleanups then status is not set :param user_list: """ users = [u["user"] for u in user_list] def expression(users) -> Tuple[str, Literal["running", "stopped", None]]: """ Returns a list of (user, cleanup_status) tuples where `cleanup_status` can take one of the values: "running", "stopped", or None """ case_running = Case( None, ( ( MalwareHit.status.in_( ( MalwareHitStatus.CLEANUP_PENDING, MalwareHitStatus.CLEANUP_STARTED, ) ), 1, ), ), 0, ) case_stopped = Case( None, ( ( MalwareHit.status.in_( ( MalwareHitStatus.CLEANUP_DONE, MalwareHitStatus.CLEANUP_REMOVED, ) ), 1, ), ), 0, ) query = ( MalwareHit.select( MalwareHit.user, Case( None, ( (fn.Sum(case_running) > 0, "running"), (fn.Sum(case_stopped) > 0, "stopped"), ), ).alias("cleanup_status"), ) .where(MalwareHit.user.in_(users)) .group_by(MalwareHit.user) ) return query.tuples() for user, status in get_results_iterable_expression(expression, users): u = get(user_list, user=user) u["cleanup_status"] = status def update_last_scan_date(user_list): if is_cluster(): # Apps share a homedir, so attribute the last scan date per app # via its prefixed scan path rather than to everyone under the # shared bare path. by_user = {u["user"]: u for u in user_list} for user, scan in _latest_completed_scan_per_app(user_list): by_user[user]["scan_date"] = scan.completed return def expression(homes): return ( MalwareScan.select(MalwareScan.path, MalwareScan.completed) .where(MalwareScan.path.in_(homes)) .group_by(MalwareScan.path) .having(MalwareScan.completed == fn.Max(MalwareScan.completed)) ) home_to_users = defaultdict(list) for user in user_list: home_to_users[user["home"]].append(user) for scan in get_results_iterable_expression( expression, list(home_to_users) ): for user in home_to_users[scan.path]: user["scan_date"] = scan.completed def update_cleanup_analyst_status(user_list): """ Updates cleanup analyst status for the list of panel users. Checks if users have active cleanup requests (pending or in_progress) or recently completed requests (within the last 3 days). :param user_list: List of user dictionaries to update """ if not user_list: return # Calculate the cutoff date for "recently completed" (3 days ago) three_days_ago = datetime.now(timezone.utc) - timedelta(days=3) # Get all usernames from the user list usernames = [u["user"] for u in user_list if u["user"]] def expression(users): """ Returns a query to fetch active cleanup requests and recently completed requests for the specified users. """ return ( AnalystCleanupRequest.select( AnalystCleanupRequest.username, AnalystCleanupRequest.status, AnalystCleanupRequest.last_updated, AnalystCleanupRequest.created_at, ) .where( (AnalystCleanupRequest.username.in_(users)) & ( ( AnalystCleanupRequest.status.in_( ["pending", "in_progress"] ) ) | ( (AnalystCleanupRequest.status == "completed") & ( AnalystCleanupRequest.last_updated >= three_days_ago ) ) ) ) .order_by(AnalystCleanupRequest.created_at.desc()) ) # Use the same get_results_iterable_expression pattern as other update functions requests = get_results_iterable_expression(expression, usernames) # Create a mapping of username to their most recent request user_to_request = {} for request in requests: # If we already have a request for this user and it's not newer, skip if ( request.username in user_to_request and user_to_request[request.username].created_at > request.created_at ): continue user_to_request[request.username] = request # Update each user with their request status for u in user_list: username = u["user"] if username in user_to_request: request = user_to_request[username] u["analyst_status"] = request.status async def get_matched_users(match) -> Tuple[int, List[Dict[str, Any]]]: user_list = await panel_users() if isinstance(match, str): pattern = re.compile(f".*{match}.*") elif isinstance(match, Iterable): pattern = re.compile(f"^({'|'.join(match)})$") else: pattern = re.compile(".*") matched_users = [u for u in user_list if pattern.match(u["user"])] return len(user_list), matched_users async def fetch_user_list(get_scans, *, match=None): max_count, user_list = await get_matched_users(match) update_infected_count_and_last_scan(user_list) update_running_scan_status(user_list, get_scans) update_cleanup_status(user_list) update_last_scan_date(user_list) if await AnalystCleanupAPI.check_cleanup_allowed(): update_cleanup_analyst_status(user_list) return max_count, user_list def sort(user_list, field="_infected_total", desc=True): def getter(element): field_type = ( int if field in ["infected", "infected_db", "_infected_total", "scan_date"] else str ) min_val = chr(0) if field_type is str else 0 value = element.get(field) if value is None: value = min_val return value user_list.sort(key=getter, reverse=desc) if field == "_infected_total": for user in user_list: user.pop("_infected_total") return user_list async def get_file_owner( path: str, users_from_panel: Set[str], pw_all: List[pwd.struct_passwd] ): """Get username, groupname for file *path* pw_all - should contains result of pwd.getpwall() user_from_panel - users in current panel (see code comment) Returns tuple (user, group, uid, gid)""" stat = os.stat(path) owner = user = uid = stat.st_uid group = gid = stat.st_gid p = Path(path) if uid == 0 and is_crontab(p): for pw in pw_all: if pw.pw_name == p.name: if pw.pw_name in users_from_panel: owner, user, uid = pw.pw_name, pw.pw_name, pw.pw_uid group = gid = pw.pw_gid break else: # Plesk-panel clients can have two system users with same uids, # but only one of them will be used in panel. for pw in pw_all: if pw.pw_uid == uid and pw.pw_name in users_from_panel: owner = user = pw.pw_name break try: group = (await to_thread(grp.getgrgid, gid)).gr_name except KeyError: pass user = CloudwaysUser.override_name_by_path( p, user, users_from_panel, pw_all ) return owner, user, group, uid, gid async def fill_results_owner(results, *, scan_id=None): """Fill owner info for scan results. In cluster mode the tenant is the path's first component, so each result is labelled from its own key: the agent cannot stat a file living in another container, and a batch spanning tenants must not be given one of them wholesale. Paths naming no tenant still go through the stat below. """ remaining = results if is_cluster(): # the split is syntactic, so a first component only names a tenant # when it is a registered application: /var/www/... is not app "var" registered = await registered_apps() remaining = {} for path, data in results.items(): tenant, _ = split_prefixed(path) if tenant and tenant in registered: data["owner"] = tenant data["user"] = tenant data["group"] = tenant data["uid"] = 0 data["gid"] = 0 else: remaining[path] = data if not remaining: return users_from_panel = set(await hosting_panel.HostingPanel().get_users()) missing = [] pw_all = await to_thread(pwd.getpwall) for path, data in remaining.items(): try: owner, user_from_stat, group, uid, gid = await get_file_owner( path, users_from_panel, pw_all, ) except FileNotFoundError: missing.append(path) else: data["owner"] = owner data["user"] = user_from_stat data["group"] = group data["uid"] = uid data["gid"] = gid for m in missing: del results[m] if missing: logger.warning( "Dropped %d scan result(s) for files gone before owner" " resolution (scan %s): %s", len(missing), scan_id, missing, ) # stable format string: one Sentry issue, event count = fleet rate logger.error( "Scan results dropped: %d file(s) vanished between scan and" " report", len(missing), extra={"scan_id": scan_id, "paths": missing}, ) def is_uid(username: str) -> bool: try: int(username) return True except ValueError: return False async def get_username_by_uid(uid) -> str: uid = int(uid) pw_all = await to_thread(pwd.getpwall) return next( (pw.pw_name for pw in pw_all if pw.pw_uid == uid), None, )