"""
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,
)