"""
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 <https://www.gnu.org/licenses/>.

Copyright © 2019 Cloud Linux Software Inc.

This software is also available under ImunifyAV commercial license,
see <https://www.imunify360.com/legal/eula>
"""
import asyncio
import errno
import json
import logging
import os
import signal as signal_module
import subprocess
import tempfile
import time
from collections import defaultdict
from contextlib import asynccontextmanager, suppress
from itertools import islice
from pathlib import Path
from typing import Callable, Dict, Generator, List, Optional, Set, Tuple, Union

from defence360agent.contracts.config import (
    Malware,
    MalwareSignatures,
    MyImunifyConfig,
)
from defence360agent.contracts.license import LicenseCLN
from defence360agent.contracts.messages import MessageType
from defence360agent.contracts.permissions import (
    ms_clean_requires_myimunify_protection,
)
from defence360agent.utils import (
    RecurringCheckStop,
    Singleton,
    base64_encode_filename,
    recurring_check,
    resource_limits,
)
from imav.contracts.config import MalwareTune
from imav.malwarelib.engine import curator_path
from imav.malwarelib.model import MalwareHit
from imav.malwarelib.utils.revisium import (
    DeletionType,
    ErrorType,
    RescanResultType,
    RevisiumCSVFile,
    RevisiumJsonFile,
    RevisiumTempFile,
)

logger = logging.getLogger(__name__)


def cleaner_result_instance(tempdir=None, mode=None):
    if MalwareTune.USE_JSON_REPORT:
        return RevisiumJsonFile(tempdir, mode)
    return RevisiumCSVFile(tempdir, mode)


def usable_tempdir() -> str:
    """The temp dir both the agent and the privilege-dropped cleaner can use."""
    tempdir = Path(tempfile.gettempdir())
    # gettempdir() falls back to the cwd when TMPDIR, /tmp, /var/tmp and
    # /usr/tmp are all unusable, and the agent unit has no WorkingDirectory,
    # so that fallback is "/" — where nothing can be created.
    if tempdir == Path(os.sep) or not os.access(tempdir, os.W_OK | os.X_OK):
        raise PermissionError(
            errno.EACCES, "no usable temp directory", str(tempdir)
        )
    return str(tempdir)


class CleanerMemory:
    SAMPLE_INTERVAL = 2

    def __init__(self, pid):
        self._pid = pid
        self._peak_rss_kb = None
        self._bound = self._read(resource_limits.memory_bound_kb)
        self.sample()

    def _read(self, probe):
        try:
            return probe(self._pid)
        except OSError:
            return None

    def sample(self):
        peak = self._read(resource_limits.peak_rss_kb)
        if peak:
            self._peak_rss_kb = max(peak, self._peak_rss_kb or 0)

    async def watch(self):
        while True:
            await asyncio.sleep(self.SAMPLE_INTERVAL)
            self.sample()

    def note(self) -> Optional[str]:
        if not self._peak_rss_kb:
            return None
        peak_mb = self._peak_rss_kb // 1024
        if self._bound:
            bound_kb, source = self._bound
            return f"peak RSS {peak_mb} MB of {bound_kb // 1024} MB, {source}"
        # Cleanup is spawned outside run-with-intensity, so on CloudLinux the
        # ceiling that killed it is the user's LVE, which we cannot read here.
        lve = ""
        with suppress(OSError):
            if resource_limits.is_lve_active():
                lve = ", LVE active"
        return f"peak RSS {peak_mb} MB, no cgroup limit{lve}"


class MalwareCleanerLog(RevisiumTempFile):
    pass


class MalwareCleanerProgress(RevisiumJsonFile):
    """
    Get progress from external source
    """

    _progress = 0

    @recurring_check(2)
    async def watch(self, callback):
        try:
            data = self.read()
        except FileNotFoundError:
            raise RecurringCheckStop()
        except json.JSONDecodeError:
            return

        progress = data["current"]

        increment, self._progress = progress - self._progress, progress

        callback(increment)


class MalwareCleanupFileList(RevisiumTempFile):
    def write(self, filelist):
        with self.open_for_write() as w:
            w.writelines(base64_encode_filename(f) + b"\n" for f in filelist)


def _normalize_path(path: str) -> str:
    """Round-trip through fsencode/fsdecode to match FilenameField behavior.

    Ensures CleanupResult keys are consistent with MalwareHit.orig_file
    which goes through FilenameField's os.fsencode/os.fsdecode cycle.
    """
    return os.fsdecode(os.fsencode(path))


def _parse_int(value: Union[str, int]) -> int:
    """Convert str|int to int, in case errors return -2
    -1 used as default value when storing CH
    """
    try:
        return int(value)
    except (TypeError, ValueError):
        return -2


class CleanupResultEntry(dict):
    def __init__(self, data: Dict[str, Union[str, int]]):
        # fields:
        # d - cleanup result
        # e - error description
        # s - signature that was triggered for the file during scan
        # f - file path (it's unexpected that f is absent)
        # r - the result of aibolit rescan after cleanup
        #
        # We shouldn't fail on parsing one record (to do not stop processing
        #  report), so we consider default values for all fields.
        super().__init__(
            d=_parse_int(data.get("d", -1)),
            e=_parse_int(data.get("e", -1)),
            s=data["s"],
            f=data["f"],
            r=_parse_int(data.get("r", -1)),
            mtime_before=_parse_int(data.get("mb", -1)),
            mtime_after=_parse_int(data.get("ma", -1)),
            hash_before=data.get("hb", ""),
            hash_after=data.get("ha", ""),
        )

    def is_cleaned(self):
        if self.is_failed() or self.requires_myimunify_protection():
            return False

        if self["e"] == ErrorType.NOT_CLEANEDUP:
            logger.warning(
                "File has changed, assuming that it was cleaned: %s", self["f"]
            )
            return True

        return (
            self["e"] == ErrorType.NO_ERROR
            and self["d"] == DeletionType.INJECTION_REMOVED
        )

    def is_removed(self):
        return (
            not self.is_failed()
            and self["e"] == ErrorType.NO_ERROR
            and self["d"] > DeletionType.INJECTION_REMOVED
        )

    def is_failed(self):
        # a system-owner refusal carries a synthetic r=1 for older agents;
        # it is not a failed cleanup attempt and must not go to MRS
        return (
            self["r"] == RescanResultType.DETECTED
            and self["e"] != ErrorType.SKIPPED_SYSTEM_OWNER
        )

    def requires_myimunify_protection(self):
        return self["r"] == RescanResultType.REQUIRED_ADVANCED_SIGNATURES

    def not_exist(self):
        return not self.is_failed() and self["e"] == ErrorType.FILE_NOT_EXISTS


class CleanupResult(Dict[str, CleanupResultEntry]):
    """
    Cleanup result container for result entries
    """

    def __init__(self, report=None):
        if not report:
            return
        named = [entry for entry in report if entry.get("f")]
        if len(named) != len(report):
            logger.warning(
                "Dropped %d cleanup report entries without a path out of %d",
                len(report) - len(named),
                len(report),
            )
        super().__init__(
            {
                _normalize_path(entry["f"]): CleanupResultEntry(entry)
                for entry in named
            }
        )

    @staticmethod
    def __key(hit: Union[str, MalwareHit]) -> str:
        return getattr(hit, "orig_file", hit)

    def __contains__(self, hit: Union[str, MalwareHit]):
        return super().__contains__(self.__key(hit))

    def __getitem__(self, hit: Union[str, MalwareHit]):
        return super().__getitem__(self.__key(hit))


class MalwareCleaner:
    PROCU_DB = MalwareSignatures.PROCU_DB

    def __init__(self, loop=None, sink=None, watch_progress=True):
        self._loop = loop if loop else asyncio.get_event_loop()
        self._proxy = MalwareCleanupProxy()
        self._sink = sink
        self._watch_progress = watch_progress

    def _cmd(
        self,
        filename,
        progress_path,
        result_path,
        log_path,
        soft,
        *,
        username,
        blacklist=None,
        use_csv=True,
        standard_only=True,
    ):
        cleaner_binary = curator_path()
        cmd = [
            "/opt/ai-bolit/wrapper",
            cleaner_binary,
            "--deobfuscate",
            "--nobackup",
            "--forcibly_cleanup",
            "--rescan",
            "--list=%s" % filename,
            "--input-fn-b64-encoded",
            "--username=%s" % username,
            "--report-hashes",
        ]
        if blacklist:
            cmd.append("--black-list=%s" % blacklist)
        cmd.extend(
            [
                "--log=%s" % log_path,
                "--progress=%s" % progress_path,
            ]
        )
        if (
            Malware.CLEANUP_DISABLE_CLOUDAV
            or not LicenseCLN.is_cloud_assisted_cleanup_allowed()
        ):
            cmd.append("--disable-cloudav")

        if use_csv:
            cmd.extend(["--csv_result=%s" % result_path])
        else:
            cmd.extend(["--result=%s" % result_path])

        if standard_only:
            cmd.extend(["--standard-only"])

        if os.path.exists(self.PROCU_DB):
            cmd.append("--avdb")
            cmd.append(self.PROCU_DB)
        if soft:
            cmd.append("--soft")

        return cmd

    CLEANUP_ERROR_PREFIX = "Cleanup failed."

    # Procu exit codes (must match constants in procu2_src.php)
    _EXIT_ERROR_GENERAL = 1
    _EXIT_ERROR_INPUT_NOT_FOUND = 2
    _EXIT_ERROR_INVALID_USERNAME = 255
    _SIGNAL_EXIT_CODE_OFFSET = 128

    @staticmethod
    def _get_cleaner_error_info(
        exc: Exception,
        cmd: List[str],
        returncode: int,
        stdout: Optional[bytes],
        stderr: Optional[bytes],
    ):
        return dict(
            exception=exc.__class__.__name__,
            return_code=returncode,
            command=cmd,
            out=stdout.decode(errors="replace") if stdout is not None else "",
            err=stderr.decode(errors="replace") if stderr is not None else "",
        )

    @classmethod
    def _categorize_process_error(
        cls,
        returncode: int,
        stderr: bytes,
        memory_note: Optional[str] = None,
    ) -> Optional[str]:
        """Categorize procu process failures into explicit error types.

        Returns a human-readable error string if the process failed,
        or None on success (exit code 0).
        """
        if returncode == 0:
            return None

        prefix = cls.CLEANUP_ERROR_PREFIX
        stderr_text = stderr.decode(errors="replace") if stderr else ""

        # asyncio reports -N for a signalled child; the wrapper reports 128+N.
        if returncode < 0:
            sig_num = -returncode
        else:
            sig_num = returncode - cls._SIGNAL_EXIT_CODE_OFFSET

        sig_name = None
        if sig_num > 0:
            with suppress(ValueError):
                sig_name = signal_module.Signals(sig_num).name

        if sig_name:
            # Only the signal categories carry the note: the other strings are
            # matched exactly by the cleanup dashboards.
            memory = f" [{memory_note}]" if memory_note else ""
            if sig_num == signal_module.SIGSEGV:
                return (
                    f"{prefix} Segmentation fault (signal {sig_name}){memory}"
                )
            if sig_num == signal_module.SIGKILL:
                return (
                    f"{prefix} Process killed (signal {sig_name},"
                    f" likely OOM){memory}"
                )
            if sig_num == signal_module.SIGTERM:
                return (
                    f"{prefix} Process terminated (signal {sig_name}){memory}"
                )
            return f"{prefix} Process killed by signal {sig_name}{memory}"

        if (
            "Allowed memory size of" in stderr_text
            and "bytes exhausted" in stderr_text
        ):
            return f"{prefix} Out of memory"
        if "Fatal error" in stderr_text:
            return f"{prefix} PHP fatal error (exit code {returncode})"
        if "Parse error" in stderr_text:
            return f"{prefix} PHP parse error (exit code {returncode})"
        if "error while loading shared libraries" in stderr_text:
            return f"{prefix} Shared library error (exit code {returncode})"

        if returncode == cls._EXIT_ERROR_GENERAL:
            return f"{prefix} General error (exit code 1)"
        if returncode == cls._EXIT_ERROR_INPUT_NOT_FOUND:
            return f"{prefix} Input file not found (exit code 2)"
        if (
            returncode == cls._EXIT_ERROR_INVALID_USERNAME
            and "Invalid username" in stderr_text
        ):
            return f"{prefix} Invalid username (exit code 255)"

        return f"{prefix} Process exited with code {returncode}"

    @classmethod
    def _read_report(cls, result) -> Tuple[Optional[list], Optional[str]]:
        """Returns (report, error); the error is None when the read worked."""
        try:
            return result.read(), None
        except FileNotFoundError:
            return None, f"{cls.CLEANUP_ERROR_PREFIX} Report file is missing"
        except Exception as exc:
            logger.warning(
                "Cleanup report %s is unreadable: %s: %s",
                result.filename,
                exc.__class__.__name__,
                exc,
            )
            error = (
                f"{cls.CLEANUP_ERROR_PREFIX} Report is unreadable"
                f" ({exc.__class__.__name__})"
            )
            return None, error

    async def _fail_before_run(
        self, exc: OSError
    ) -> Tuple[CleanupResult, Optional[str], List[str]]:
        error = (
            f"{self.CLEANUP_ERROR_PREFIX} No usable temp directory"
            f" ({exc.__class__.__name__})"
        )
        logger.error("%s: %s", error, exc)
        info = self._get_cleaner_error_info(
            exc, [], -1, stdout=b"", stderr=b""
        )
        await self._send_cleanup_failed_message({**info, "message": error})
        return CleanupResult(), error, []

    @asynccontextmanager
    async def _watch_memory(self, pid):
        memory = CleanerMemory(pid)
        watcher = self._loop.create_task(memory.watch())
        try:
            yield memory
        finally:
            watcher.cancel()
            with suppress(asyncio.CancelledError):
                await watcher

    async def _send_cleanup_failed_message(self, info: dict):
        if self._sink:
            try:
                msg = MessageType.CleanupFailed(
                    {**info, **{"timestamp": int(time.time())}}
                )
                await self._sink.process_message(msg)
            except asyncio.CancelledError:
                raise
            except Exception:
                logger.exception(
                    "Exception while sending CleanupFailed message"
                )

    async def start(
        self,
        user,
        filelist,
        soft=True,
        blacklist=None,
        standard_only=None,
    ) -> Tuple[CleanupResult, Optional[str], List[str]]:
        standard_only = self.is_standard_only(user, standard_only)

        try:
            tempdir = usable_tempdir()
        except OSError as exc:
            return await self._fail_before_run(exc)

        result_file = cleaner_result_instance(tempdir=tempdir)
        use_csv = isinstance(result_file, RevisiumCSVFile)

        with (
            MalwareCleanupFileList(tempdir=tempdir, mode=0o644) as flist,
            MalwareCleanupFileList(tempdir=tempdir, mode=0o644) as blk,
            MalwareCleanerProgress(tempdir=tempdir) as progress,
            result_file as result,
            MalwareCleanerLog(tempdir=tempdir) as log,
        ):
            flist.write(filelist)
            if blacklist:
                blk.write(blacklist)
            if self._watch_progress:
                self._loop.create_task(progress.watch(self._proxy.progress_cb))

            cmd = self._cmd(
                flist.filename,
                progress.filename,
                result.filename,
                log.filename,
                soft,
                username=user,
                blacklist=blk.filename if blacklist else None,
                use_csv=use_csv,
                standard_only=standard_only,
            )
            logger.debug("Executing %s", " ".join(cmd))

            out, err = b"", b""
            proc = None
            memory = None
            try:
                proc = await asyncio.subprocess.create_subprocess_exec(
                    *cmd,
                    stdout=subprocess.PIPE,
                    stderr=subprocess.PIPE,
                )
                async with self._watch_memory(proc.pid) as memory:
                    out, err = await proc.communicate()
            except asyncio.CancelledError:
                if proc:
                    with suppress(ProcessLookupError):
                        proc.terminate()
                raise
            except Exception as exc:
                info = self._get_cleaner_error_info(
                    exc,
                    cmd,
                    proc.returncode if proc else 126,  # 126 - permission error
                    stdout=out,
                    stderr=err,
                )
                # Exit code in the template, so Sentry groups by it
                logger.error(
                    f"Cleanup failed exit_code={info['return_code']}: %s",
                    f"{info['out']} {info['err']}",
                    extra={**info, "exception": exc},
                )
                await self._send_cleanup_failed_message(
                    {**info, "message": str(exc)}
                )
                return (
                    CleanupResult(),
                    (
                        f"{self.CLEANUP_ERROR_PREFIX} Failed to run the"
                        f" cleaner ({exc.__class__.__name__})"
                    ),
                    cmd,
                )

            # Categorize before reading, or the exit code is lost to the read.
            error = self._categorize_process_error(
                proc.returncode, err, memory.note() if memory else None
            )
            report, read_error = self._read_report(result)
            error = error or read_error
            if not error and not report:
                error = f"{self.CLEANUP_ERROR_PREFIX} Report is empty"

            if error:
                logger.error(
                    "%s (exit_code=%d, input_files=%d, blacklisted=%d,"
                    " stderr=%s)",
                    error,
                    proc.returncode,
                    len(filelist),
                    len(blacklist or ()),
                    err,
                )
                info = self._get_cleaner_error_info(
                    RuntimeError(error),
                    cmd,
                    proc.returncode,
                    stdout=out,
                    stderr=err,
                )
                await self._send_cleanup_failed_message(
                    {**info, "message": error}
                )
            elif len(report) < len(filelist):
                logger.warning(
                    "Partial cleanup report: %d entries for %d input files",
                    len(report),
                    len(filelist),
                )

            return CleanupResult(report), error, cmd

    @staticmethod
    def is_standard_only(user: str, standard_only: bool) -> bool:
        """Check if only standard signatures should be applied for the user"""

        # FIXME: DEF-20763 Remove this line to enable standard signatures
        return False

        if not MyImunifyConfig.ENABLED:
            # Ignore standard_only value if MyImunify is disabled
            return False
        elif standard_only is None:
            # When cleaned by default action
            return not ms_clean_requires_myimunify_protection(user)

        return standard_only


class MalwareCleanupProxy(metaclass=Singleton):
    _CHUNK_SIZE = 10000
    """
    Class to interconnect Cleanup status endpoint and Cleanup plugin
    """

    def __init__(self):
        self.current = self.total = 0
        self.hits = defaultdict(set)

    def add(self, cause, initiator, post_action, scan_id, standard_only, hits):
        self.hits[
            (cause, initiator, post_action, scan_id, standard_only)
        ].update(hits)

    def flush(
        self,
    ) -> Generator[Tuple[str, str, Callable, str, Set], None, None]:
        while self.hits:
            scan_info, hits = self.hits.popitem()

            all_hits = iter(hits)
            hits = set(islice(all_hits, self._CHUNK_SIZE))

            remaining_hit = next(all_hits, None)
            if remaining_hit is not None:
                self.hits[scan_info].add(remaining_hit)
                self.hits[scan_info].update(all_hits)

            self.total += len(hits)
            yield *scan_info, hits

    def progress_cb(self, increment=1):
        self.current += increment

    def reset(self):
        self.current = self.total = 0

    def get_progress(self):
        try:
            return int(self.current / (self.total + len(self.hits)) * 100)
        except ZeroDivisionError:
            return None
