import asyncio
import logging
import pwd
from contextlib import suppress
from pathlib import Path
from typing import Coroutine

from peewee import OperationalError

from defence360agent.contracts.config import (
    ANTIVIRUS_MODE,
    ConfigValidationError,
    SystemConfig,
    UserConfig,
    Wordpress,
)
from defence360agent.contracts.hook_events import HookEvent
from defence360agent.contracts.license import LicenseCLN
from defence360agent.contracts.messages import MessageType
from defence360agent.contracts.plugins import (
    MessageSink,
    MessageSource,
    expect,
)
from defence360agent.internals.components import WAF_SITES_STATE
from defence360agent.subsys.panels import hosting_panel
from defence360agent.subsys.persistent_state import (
    load_state,
    register_lock_file,
    save_state,
)
from defence360agent.utils import (
    Scope,
    importer,
    recurring_check,
    system_packages_info,
)
from defence360agent.utils.check_lock import check_lock
from defence360agent.internals.iaid import IndependentAgentIDAPI
from defence360agent.utils.common import DAY, HOUR

from defence360agent.wordpress import cli as wp_cli
from defence360agent.wordpress import plugin
from defence360agent.wordpress.utils import _prepare_ai_bot_settings

from defence360agent.wordpress.bot_protection import (
    resolve_ai_bot_protection,
)
from defence360agent.model import tls_check
from defence360agent.model.wordpress import WPSite, WordpressSite
from defence360agent.wordpress.site_repository import (
    get_sites_by_path,
    get_sites_for_user,
    get_installed_sites,
)
from defence360agent.wordpress.proxy_auth import (
    is_secret_expired,
    rotate_secret,
)
from defence360agent.wordpress import (
    BotStatsCollector,
    ChangelogProcessor,
    IncidentCollector,
    IncidentSender,
)
from defence360agent.wordpress.plugin import update_disabled_rules_on_sites
from defence360agent.model.wordpress_incident import (
    delete_old_wordpress_incidents,
)


logger = logging.getLogger(__name__)

LOCK_FILE = register_lock_file("wp-gen-auth", Scope.AV_IM360)
SITE_PROCESSING_LOCK_FILE = register_lock_file(
    "wp-site-process", Scope.AV_IM360
)
SEND_WP_PLUGIN_STATS_LOCK_FILE = register_lock_file(
    "wp-plugin-stats", Scope.AV_IM360
)
SEND_WP_PLUGIN_BOT_STATS_LOCK_FILE = register_lock_file(
    "wp-plugin-bot-stats", Scope.AV_IM360
)
LICENSE_RECONVERGE_LOCK_FILE = register_lock_file(
    "wp-license-reconverge", Scope.AV_IM360
)
IP_WHITELIST_LOCK_FILE = register_lock_file("wp-ip-whitelist", Scope.AV_IM360)

CONFIG_DIR = Path("/etc/sysconfig/imunify360/imunify360.config.d")

FIRST_INSTALL_CONFIG_FILE = Path(
    "/opt/imunify360/venv/share/imunify360/11_on_first_install_wp_av.config"
)
FIRST_INSTALL_CONFIG_PATH = CONFIG_DIR / "11_on_first_install_wp_av.config"

FIRST_INSTALL_FLAG = CONFIG_DIR / ".11_on_first_install_wp_av.flag"

_MalwareHit = importer.get(
    module="imav.malwarelib.model", name="MalwareHit", default=None
)


def _get_cleaned_malware_hits(started_timestamp: float) -> list:
    """
    Get malware hits cleaned since the given timestamp with lazy import fallback.

    Returns empty list if imav.malwarelib is not available.
    """
    if _MalwareHit is None:
        logger.debug(
            "imav.malwarelib not available, returning empty cleaned hits"
        )
        return []
    return _MalwareHit.cleaned_since(started_timestamp)


class ImunifySecurityPlugin(MessageSink, MessageSource):
    SCOPE = Scope.AV_IM360

    def __init__(self):
        self._loop = None
        self._sink = None
        state = load_state("ImunifySecurityPlugin")
        self.installation_completed = state.get("installed")
        # Explicit None check: a persisted False must win over the live
        # config value. Plain `or` would flip False → True via short-circuit.
        persisted_enabled = state.get("enabled")
        self.last_config_value = (
            persisted_enabled
            if persisted_enabled is not None
            else Wordpress.SECURITY_PLUGIN_ENABLED
        )
        self._last_waf_enabled = plugin._get_global_waf_enabled()
        self._last_waf_default = plugin._get_waf_default()
        self._last_user_waf_enabled: dict[str, bool | None] = {}
        # Seed with the current value so the first ConfigUpdate after startup
        # only triggers a propagation when the admin has actually toggled it.
        self._last_ai_bot_protection = plugin._get_global_ai_bot_protection()
        self._last_ai_bot_protection_preset = (
            plugin._get_global_ai_bot_protection_preset()
        )
        # Seed the RAW nullable value so it compares like-for-like with the
        # reactor's raw read; a resolved bool here would look "changed"
        # against an unset (None) default and fire a needless redeploy.
        try:
            self._last_ai_bot_protection_default = (
                Wordpress.AI_BOT_PROTECTION_DEFAULT
            )
        except KeyError:
            self._last_ai_bot_protection_default = None
        self._last_user_ai_bot_protection: dict[str, tuple] = {}
        # Seeded None (not a live read) to keep license-file I/O out of the
        # constructor; the poll reconverges on the first tick if it differs.
        self._last_license_type = None
        # Rendered ip_whitelist.php text last dispatched to the sites. A
        # dispatch marker, not an apply receipt: None until the first pass
        # after startup/(re)install so that pass runs against every site
        # (atomic_rewrite makes it read-only where the file is current).
        self._last_ip_whitelist: str | None = None
        self._ip_whitelist_lock = asyncio.Lock()
        self._ip_whitelist_size_warned = False
        self.installation_task: asyncio.Task | None = None
        self.deleting_task: asyncio.Task | None = None
        self.install_and_update_task: asyncio.Task | None = None
        self.freshly_installed_sites: set[WPSite] = set()

        # Incident collection and changelog processing components
        self.incident_collector = IncidentCollector()
        self.bot_stats_collector = BotStatsCollector()
        self.incident_sender = IncidentSender()
        self.changelog_processor = ChangelogProcessor()
        self._site_processing_task: asyncio.Task | None = None
        self._stats_task: asyncio.Task | None = None
        self._bot_stats_task: asyncio.Task | None = None
        self._license_reconverge_task: asyncio.Task | None = None
        self._ip_whitelist_task: asyncio.Task | None = None
        self._ai_bot_reconcile_task: asyncio.Task | None = None

    async def create_sink(self, loop):
        pass

    async def create_source(self, loop, sink):
        self._loop = loop
        self._sink = sink
        self._update_auth_task = self._loop.create_task(
            self.refresh_auth_files()
        )

        self._site_processing_task = self._loop.create_task(
            self.process_wordpress_sites()
        )
        self._stats_task = self._loop.create_task(self.send_stats())
        self._bot_stats_task = self._loop.create_task(self.send_bot_stats())
        self._license_reconverge_task = self._loop.create_task(
            self.reconverge_license_type()
        )
        self._ip_whitelist_task = self._loop.create_task(
            self.refresh_ip_whitelist()
        )

        if ANTIVIRUS_MODE:
            await self._apply_first_install_config()
        else:
            FIRST_INSTALL_FLAG.unlink(missing_ok=True)

        await self._recover_installation_on_startup()

    async def _recover_installation_on_startup(self):
        """
        Self-heal when the installation state was lost.

        If the feature is enabled but installation_completed is falsy (state
        file missing, earlier install interrupted, etc.), manage_plugin_installation
        can never recover: its True == True guard always returns early.
        Trigger install_everywhere once per restart to repopulate the
        wordpress_site table and flip the flag.
        """
        if (
            self.installation_completed
            or not Wordpress.SECURITY_PLUGIN_ENABLED
        ):
            return
        logger.info(
            "Installation state is missing while feature is enabled; "
            "triggering startup self-recovery"
        )
        # Sync last_config_value to the live config. If the persisted state
        # held a stale `enabled: False`, _mark_installation_done would bail
        # on `if not self.last_config_value` and installation_completed would
        # never flip — recovery would re-run on every restart.
        self.last_config_value = True
        await self.process_installation(
            plugin.install_everywhere(sink=self._sink)
        )
        if self.installation_task is not None:
            self.installation_task.add_done_callback(
                self._mark_installation_done
            )

    async def _apply_first_install_config(self):
        if not FIRST_INSTALL_FLAG.exists():
            return
        if await hosting_panel.HostingPanel().users_count() == 1:
            _ = FIRST_INSTALL_CONFIG_PATH.write_text(
                FIRST_INSTALL_CONFIG_FILE.read_text()
            )
            FIRST_INSTALL_CONFIG_PATH.chmod(0o600)
        FIRST_INSTALL_FLAG.unlink()

    async def shutdown(self):
        self._update_auth_task.cancel()
        # CancelledError is handled by @recurring_check():
        await self._update_auth_task

        # Cancel site processing (changelogs + incidents) task
        if self._site_processing_task:
            self._site_processing_task.cancel()
            await self._site_processing_task

        if self._stats_task:
            self._stats_task.cancel()
            await self._stats_task

        if self._bot_stats_task:
            self._bot_stats_task.cancel()
            await self._bot_stats_task

        if self._license_reconverge_task:
            self._license_reconverge_task.cancel()
            await self._license_reconverge_task

        if self._ip_whitelist_task:
            self._ip_whitelist_task.cancel()
            await self._ip_whitelist_task

        await self._cancel_ai_bot_reconcile()

    def _task_in_progress(self, task_attr_name):
        if not hasattr(self, task_attr_name):
            logger.error("Unknown task '%s'", task_attr_name)
            return False
        task = getattr(self, task_attr_name)

        return task is not None and not task.done() and not task.cancelled()

    def _save_installation_state(self):
        save_state(
            "ImunifySecurityPlugin",
            {
                "installed": self.installation_completed,
                "enabled": self.last_config_value,
            },
        )

    def _mark_installation_done(self, task: asyncio.Task):
        if task.cancelled():
            logger.info("Installation task was cancelled")
            return

        exc = task.exception()
        if exc is not None:
            logger.error("Installation task failed: %s", exc)
            return

        if not self.last_config_value:
            logger.info(
                "Feature was disabled during installation, "
                "skipping flag update"
            )
            return

        self.installation_completed = True
        self._save_installation_state()
        # The install pass rewrote the data directories, so the cached
        # export no longer reflects disk; the next tick writes every site.
        # This runs outside the export lock, so a poll already in flight can
        # re-cache the current text and swallow the reset — harmless only
        # because the installer writes each site's file itself; do not rely
        # on this reset as the sole convergence trigger.
        self._last_ip_whitelist = None
        if self._loop is not None:
            # Hold a reference: the loop keeps only a weak one, so an unstored
            # task can be garbage-collected mid-run.
            self._ai_bot_reconcile_task = self._loop.create_task(
                self._reconcile_ai_bot_after_install()
            )

    async def _reconcile_ai_bot_after_install(self):
        """Rewrite ai-bot config on every site once install completes.

        manage_ai_bot_protection_config skips per-account propagation while
        installation_completed is False, so an ai-bot value co-submitted with
        (or changed during) the enable can land wrong in plugin_config.php.
        With the installer finished there is no writer to race, so reconcile
        the file to the settled config once.
        """
        if (
            not Wordpress.SECURITY_PLUGIN_ENABLED
            or not self.installation_completed
        ):
            # Disabled or (re)installing since this task was scheduled — bail so
            # we don't rewrite files an uninstall is concurrently removing.
            return
        try:
            sites = await self._loop.run_in_executor(None, get_installed_sites)
            if sites:
                await plugin.update_plugin_config_on_sites(
                    sites, sink=self._sink
                )
        except Exception as error:
            logger.error("AI-bot post-install reconcile failed: %s", error)

    async def _cancel_ai_bot_reconcile(self):
        task = self._ai_bot_reconcile_task
        if task is not None and not task.done():
            task.cancel()
            try:
                await task
            except asyncio.CancelledError:
                pass

    async def process_installation(self, coro: Coroutine, for_new_sites=False):
        if not for_new_sites:
            # The new-sites path installs only on not-yet-installed sites,
            # disjoint from the reconcile's existing-site rewrite; cancelling
            # here would drop a just-scheduled post-install reconcile.
            await self._cancel_ai_bot_reconcile()
        if self._task_in_progress("deleting_task"):
            if for_new_sites:
                coro.close()
                return

            if self.deleting_task:
                self.deleting_task.cancel()
                try:
                    await self.deleting_task
                except asyncio.CancelledError:
                    pass

        if self._task_in_progress("installation_task"):
            logger.warning("Installation is already running")
            coro.close()
            return

        self.installation_task = asyncio.create_task(coro)

    async def process_deleting(self, coro):
        await self._cancel_ai_bot_reconcile()
        if self._task_in_progress("installation_task"):
            if self.installation_task:
                self.installation_task.cancel()
                try:
                    await self.installation_task
                except asyncio.CancelledError:
                    pass

        if self._task_in_progress("deleting_task"):
            logger.warning("Deleting is already running")
            return

        self.deleting_task = asyncio.create_task(coro)

    @recurring_check(
        check_lock,
        check_period_first=True,
        check_lock_period=DAY,
        lock_file=LOCK_FILE,
    )
    async def refresh_auth_files(self):
        if is_secret_expired():
            await rotate_secret()
        await plugin.update_auth_everywhere(sink=self._sink)

    @recurring_check(
        check_lock,
        check_period_first=True,
        jitter=True,
        check_lock_period=DAY,
        lock_file=SEND_WP_PLUGIN_STATS_LOCK_FILE,
    )
    async def send_stats(self):
        """Send WP plugin adoption stats to the correlation server."""
        if not Wordpress.SECURITY_PLUGIN_ENABLED:
            return

        sites = await self._loop.run_in_executor(None, get_installed_sites)

        def _count_manually_removed():
            with suppress(tls_check.OverridingReset):
                tls_check.reset()
            return (
                WordpressSite.select()
                .where(WordpressSite.manually_deleted_at.is_null(False))
                .count()
            )

        manually_removed = await self._loop.run_in_executor(
            None, _count_manually_removed
        )
        waf_enabled_sites = 0
        ai_bot_enabled_sites = 0
        preset_counts = {"balanced": 0, "strict": 0, "monitor": 0}
        user_ai_config: dict[int, dict] = {}

        for site in sites:
            content_dir = await wp_cli.get_content_dir(site)
            data_dir = content_dir / "imunify-security"
            rules_php = data_dir / "rules.php"
            if await self._loop.run_in_executor(None, rules_php.exists):
                waf_enabled_sites += 1

            if site.uid not in user_ai_config:
                try:
                    pw_record = await self._loop.run_in_executor(
                        None, pwd.getpwuid, site.uid
                    )
                    user_ai_config[
                        site.uid
                    ] = await self._loop.run_in_executor(
                        None,
                        _prepare_ai_bot_settings,
                        pw_record.pw_name,
                    )
                except Exception as exc:
                    logger.info(
                        "Could not load AI bot protection config for uid"
                        " %s, counting it as disabled: %s",
                        site.uid,
                        exc,
                    )
                    user_ai_config[site.uid] = {
                        "ai_bot_protection": False,
                        "preset": "balanced",
                    }

            hoster_cfg = user_ai_config[site.uid]
            ai_enabled, preset = await self._loop.run_in_executor(
                None,
                resolve_ai_bot_protection,
                site.docroot,
                data_dir,
                site.uid,
                bool(hoster_cfg.get("ai_bot_protection")),
                hoster_cfg.get("preset", "balanced"),
            )
            if ai_enabled:
                ai_bot_enabled_sites += 1
                if preset in preset_counts:
                    preset_counts[preset] += 1

        save_state(WAF_SITES_STATE, {"sites": waf_enabled_sites})

        pkgs = await system_packages_info(
            {
                "imunify360-firewall",
                "imunify-antivirus",
                "imunify-core",
                "imunify-wp-security",
            }
        )
        av_version = pkgs.get("imunify-antivirus") or ""
        firewall_version = pkgs.get("imunify360-firewall") or ""
        core_version = pkgs.get("imunify-core") or ""
        wp_version = pkgs.get("imunify-wp-security") or ""

        msg = MessageType.WpSecurityPluginStats(
            core_version=core_version,
            av_version=av_version,
            firewall_version=firewall_version,
            wp_version=wp_version,
            installed_sites=len(sites),
            manually_removed_sites=manually_removed,
            server_config={
                "waf_enabled": str(plugin._get_global_waf_enabled()),
                "ai_bot_protection": str(
                    plugin._get_global_ai_bot_protection()
                ),
            },
            stats={
                "waf_enabled_sites": str(waf_enabled_sites),
                "ai_bot_protection_enabled_sites": str(ai_bot_enabled_sites),
                "ai_bot_protection_preset_balanced": str(
                    preset_counts["balanced"]
                ),
                "ai_bot_protection_preset_strict": str(
                    preset_counts["strict"]
                ),
                "ai_bot_protection_preset_monitor": str(
                    preset_counts["monitor"]
                ),
            },
        )
        msg["iaid"] = await self._loop.run_in_executor(
            None, IndependentAgentIDAPI.get_iaid
        )
        await self._sink.process_message(msg)

    @recurring_check(
        check_lock,
        check_period_first=True,
        jitter=True,
        check_lock_period=HOUR,
        lock_file=SEND_WP_PLUGIN_BOT_STATS_LOCK_FILE,
    )
    async def send_bot_stats(self):
        """Forward each site's daily bot-traffic export to correlation."""
        # ai_bot_protection can only be turned down from the global value, so
        # with either switch off no site is producing an export to collect.
        if not (
            plugin._get_security_plugin_enabled()
            and plugin._get_global_ai_bot_protection()
        ):
            logger.debug(
                "WordPress plugin or AI bot protection is off, not looking"
                " for bot-traffic exports"
            )
            return

        try:
            sites = await self._loop.run_in_executor(None, get_installed_sites)
            await self.bot_stats_collector.collect_and_send(sites, self._sink)
        except Exception as e:
            # recurring_check abandons the loop after repeated errors, which
            # would silently stop collection until the agent restarts
            logger.warning("Bot-stats collection failed: %s", e)

    @recurring_check(
        check_lock,
        check_period_first=True,
        check_lock_period=1 * 60,  # Run every 1 minute
        lock_file=SITE_PROCESSING_LOCK_FILE,
    )
    async def process_wordpress_sites(self):
        """
        Periodic task for WordPress site file processing.

        Runs every minute to:
        1. Process changelog.php files written by the WordPress plugin (rule disable/enable from WP admin)
        2. Collect incident files written by the WordPress plugin
        """
        logger.debug(
            "Processing rule disable changelogs"
            " and collecting WordPress CVE protection incidents"
        )
        try:
            await self._process_installed_sites()
        except Exception as e:
            # deliberately not fatal to the pass below: that pass is the
            # repair path for incidents an earlier cycle failed to deliver,
            # and a collection that keeps failing must not strand them
            logger.error("Error in WordPress site processing: %s", e)

        try:
            # sends the freshly collected incidents together with any whose
            # earlier message the transport never acknowledged. Reads the
            # database, not the filesystem, so it must run even with no sites
            # left: otherwise removing the last site strands whatever the
            # transport had not yet confirmed.
            await self.incident_sender.send_pending_incidents(self._sink)
        except Exception as e:
            logger.error("Error sending pending WordPress incidents: %s", e)

    async def _process_installed_sites(self):
        sites = get_installed_sites()
        if not sites:
            logger.debug("No WordPress sites found for periodic processing")
            return

        # Process changelogs (rule disable/enable from WordPress admin)
        affected_sites = (
            await self.changelog_processor.process_changelogs_for_sites(
                sites, self._sink
            )
        )
        if affected_sites:
            await update_disabled_rules_on_sites(
                domains=[s.domain for s in affected_sites],
                sink=self._sink,
            )

        # Collect incidents
        await self.incident_collector.collect_incidents_for_sites(
            sites,
            delete_after_processing=True,
        )

        delete_old_wordpress_incidents(days=30)

    async def _install_on_new_sites(self):
        """Install plugin on new WordPress sites."""
        # Clear any previously tracked sites
        self.freshly_installed_sites.clear()

        async def install_and_track():
            installed_sites = await plugin.install_everywhere(sink=self._sink)
            if installed_sites:
                self.freshly_installed_sites.update(installed_sites)
            return installed_sites

        await self.process_installation(
            install_and_track(),
            for_new_sites=True,
        )

    async def _tidy_up(self):
        """Tidy up sites from which the WordPress plugin was deleted manually."""
        await plugin.tidy_up_manually_deleted(
            sink=self._sink,
            freshly_installed_sites=self.freshly_installed_sites,
        )
        await plugin.fix_data_file_permissions_everywhere(sink=self._sink)

        if not Wordpress.SECURITY_PLUGIN_ENABLED:
            await self.process_deleting(
                plugin.remove_all_installed(sink=self._sink)
            )
            self.installation_completed = False
            self.last_config_value = False
            self._save_installation_state()
            save_state(WAF_SITES_STATE, {"sites": 0})

    async def _adopt_found_sites(self):
        """Adopt sites where plugin is installed but not tracked in our database."""
        adopted_sites = await plugin.adopt_found_sites(sink=self._sink)
        # Add adopted sites to freshly_installed_sites to prevent them from being
        # marked as manually deleted by tidy_up (AVD database may not be updated yet)
        if adopted_sites:
            self.freshly_installed_sites.update(adopted_sites)

    async def _update_existing(self):
        """Update plugin on all sites where it is installed."""
        await plugin.update_everywhere(sink=self._sink)

    async def _remove_waf_rules_for_waf_disabled_users(self):
        """Remove WAF files from sites whose owner is WAF-off."""
        await plugin.remove_waf_rules_for_waf_disabled_users()

    async def _run_install_and_update(self):
        """
        Combined operation: install on new sites, adopt found sites, tidy up,
        update existing plugins, and remove WAF files from WAF-off accounts.
        This runs all operations sequentially to avoid race conditions.
        """
        # Install plugin on new sites.
        await self._install_on_new_sites()
        # Wait for installation to complete before proceeding.
        if self.installation_task:
            await self.installation_task

        # Adopt sites where plugin is installed but not in our database.
        await self._adopt_found_sites()

        # Tidy up and update.
        await self._tidy_up()
        await self._update_existing()

        await self._remove_waf_rules_for_waf_disabled_users()

    @expect(MessageType.WordpressPluginAction)
    async def manage_plugin_action(self, message):
        logger.info(
            "ImunifySecurityPlugin received message action: %s method: %s",
            message.action,
            message.method,
        )

        # Check if install_and_update is running - it blocks all other actions
        if self._task_in_progress("install_and_update_task"):
            logger.warning(
                "Install-and-update is still running, skipping action %s",
                message.action,
            )
            return

        if message.action == "install_on_new_sites":
            if not self.installation_completed:
                # The installation is not completed yet. We cannot know reliably which sites are new.
                return

            await self._install_on_new_sites()
            return

        if self._task_in_progress("installation_task"):
            logger.warning(
                "Installation is still running, skipping action %s",
                message.action,
            )
            return

        if self._task_in_progress("deleting_task"):
            logger.warning(
                "Uninstallation is already running, skipping action %s",
                message.action,
            )
            return

        if message.action == "update_existing":
            await self._update_existing()
            return

        if message.action == "tidy_up":
            await self._tidy_up()
            return

        if message.action == "install_and_update":
            if not self.installation_completed:
                # The installation is not completed yet. We cannot know reliably which sites are new.
                logger.warning(
                    "Installation is not completed yet, skipping"
                    " install_and_update"
                )
                return

            # Run install_and_update as a background task to prevent blocking.
            # Note: No need to check if already running - the check at the top of this function
            # (line 182) already handles that case.
            self.install_and_update_task = asyncio.create_task(
                self._run_install_and_update()
            )

    @expect(MessageType.ConfigUpdate)
    async def manage_plugin_installation(self, message):
        if not isinstance(message["conf"], SystemConfig):
            return

        current_config_value = Wordpress.SECURITY_PLUGIN_ENABLED
        if current_config_value == self.last_config_value:
            return

        # Update last config value immediately to prevent multiple installations
        self.last_config_value = current_config_value

        if current_config_value and not self.installation_completed:
            # On re-enable, force waf_enabled back on and ai_bot_protection
            # back off so a value that went stale while the plugin was off
            # cannot take effect silently. Skip a key the operator set in this
            # same update — overwriting it here would discard their explicit
            # choice with no warning. The file-poll path (config_watcher)
            # carries no delta, so it still resets both.
            submitted_wp = (message.get("submitted") or {}).get(
                "WORDPRESS", {}
            )
            if "waf_enabled" not in submitted_wp:
                try:
                    SystemConfig().dict_to_config(
                        {"WORDPRESS": {"waf_enabled": True}}
                    )
                    self._last_waf_enabled = True
                except ConfigValidationError:
                    logger.debug(
                        "waf_enabled config reset skipped,"
                        " field not in schema yet"
                    )
            else:
                # Record the operator's value so manage_waf_config, which runs
                # next on this same message, sees no change and does not start
                # an all-sites WAF removal/redeploy that races the installer
                # started below.
                self._last_waf_enabled = plugin._get_global_waf_enabled()
            if "ai_bot_protection" not in submitted_wp:
                try:
                    SystemConfig().dict_to_config(
                        {"WORDPRESS": {"ai_bot_protection": False}}
                    )
                    self._last_ai_bot_protection = False
                    self._last_ai_bot_protection_preset = (
                        plugin._get_global_ai_bot_protection_preset()
                    )
                except ConfigValidationError:
                    logger.debug(
                        "ai_bot_protection config reset skipped,"
                        " field not in schema yet"
                    )
            else:
                self._last_ai_bot_protection = (
                    plugin._get_global_ai_bot_protection()
                )
                self._last_ai_bot_protection_preset = (
                    plugin._get_global_ai_bot_protection_preset()
                )
            await self.process_installation(
                plugin.install_everywhere(sink=self._sink)
            )
            if self.installation_task is not None:
                self.installation_task.add_done_callback(
                    self._mark_installation_done
                )

        elif not current_config_value and (
            self.installation_completed
            or self._task_in_progress("installation_task")
        ):
            await self.process_deleting(
                plugin.remove_all_installed(sink=self._sink)
            )
            self.installation_completed = False
            self._save_installation_state()
            save_state(WAF_SITES_STATE, {"sites": 0})

    @expect(MessageType.ConfigUpdate)
    async def manage_ai_bot_protection_config(self, message):
        """Propagate WORDPRESS.ai_bot_protection[_preset/_default] changes to
        plugin_config.php immediately, so the WP plugin picks them up at the
        next request rather than waiting for a scan cycle.

        Caches are dispatch markers, not apply receipts — a failed site write
        is not retried on the next identical ConfigUpdate; the post-install
        reconcile and scan-cycle rewrites are the recovery net."""
        if not Wordpress.SECURITY_PLUGIN_ENABLED:
            # manage_plugin_installation clears any stale value on the
            # next plugin re-enable, so nothing to write here.
            return
        if not self.installation_completed:
            # Plugin is being (re-)installed. manage_plugin_installation has
            # already settled ai_bot_protection (reset to False, or left as
            # the operator submitted it) and the install flow writes
            # plugin_config.php for every site, so propagating here would race
            # with it and leave stale values on sites the installer skips
            # (e.g. DB rows surviving a crash during a prior disable).
            return

        conf = message["conf"]
        if isinstance(conf, UserConfig):
            await self._propagate_ai_bot_user_change(conf)
            return
        if not isinstance(conf, SystemConfig):
            return

        current_enabled = plugin._get_global_ai_bot_protection()
        current_preset = plugin._get_global_ai_bot_protection_preset()
        if (
            current_enabled != self._last_ai_bot_protection
            or current_preset != self._last_ai_bot_protection_preset
        ):
            was_enabled = self._last_ai_bot_protection
            self._last_ai_bot_protection = current_enabled
            self._last_ai_bot_protection_preset = current_preset
            sites = await self._loop.run_in_executor(None, get_installed_sites)
            if sites:
                await plugin.update_plugin_config_on_sites(
                    sites, sink=self._sink
                )
            if current_enabled:
                if not was_enabled:
                    # Sites adopted while the gate was off have no export,
                    # and the text-keyed marker cannot tell: walk every
                    # site (read-only where the file is already current).
                    self._last_ip_whitelist = None
                # The whitelist export is gated on this toggle; bring it up
                # to date now instead of at the next poll tick.
                await self._refresh_ip_whitelist_once()

        try:
            current_default = Wordpress.AI_BOT_PROTECTION_DEFAULT
        except KeyError:
            pass
        else:
            if current_default != self._last_ai_bot_protection_default:
                self._last_ai_bot_protection_default = current_default
                await plugin.apply_ai_bot_default_change(sink=self._sink)

    async def _propagate_ai_bot_user_change(self, conf) -> None:
        wp = conf.config_to_dict().get("WORDPRESS", {})
        override = (
            wp.get("ai_bot_protection"),
            wp.get("ai_bot_protection_preset"),
        )
        prev = self._last_user_ai_bot_protection.get(conf.username)
        if override == (None, None):
            # No explicit ai-bot override in this write. Act only if the
            # account had a tracked override that is now cleared — otherwise
            # this is an unrelated user-config change and rewriting every one
            # of the account's sites would be spurious (mirrors the
            # waf_value-is-None handling in manage_waf_config).
            if prev is None:
                return
            self._last_user_ai_bot_protection.pop(conf.username, None)
        elif override == prev:
            return
        else:
            self._last_user_ai_bot_protection[conf.username] = override
        await plugin.redeploy_ai_bot_for_user(conf.username, sink=self._sink)

    @recurring_check(
        check_lock,
        check_period_first=True,
        check_lock_period=1 * 60,
        lock_file=LICENSE_RECONVERGE_LOCK_FILE,
    )
    async def reconverge_license_type(self):
        """Poll for license-edition changes and reconverge plugin_config.php.

        Edition changes reach only the external hook framework, never the
        message bus, so a poll is the propagation trigger.
        """
        await self._reconverge_license_type_once()

    async def _reconverge_license_type_once(self):
        if not Wordpress.SECURITY_PLUGIN_ENABLED:
            return
        if not self.installation_completed:
            return
        current = LicenseCLN.get_license_type()
        if current == self._last_license_type:
            return
        sites = await self._loop.run_in_executor(None, get_installed_sites)
        if not sites:
            self._last_license_type = current
            return
        written = await plugin.update_plugin_config_on_sites(
            sites, sink=self._sink
        )
        if written == len(sites):
            self._last_license_type = current

    @recurring_check(
        check_lock,
        check_period_first=True,
        check_lock_period=1 * 60,
        lock_file=IP_WHITELIST_LOCK_FILE,
    )
    async def refresh_ip_whitelist(self):
        """Poll the manual IP whitelist and export it to the sites.

        Whitelist changes go CLI → im360 RPC → Go resident and never reach
        the message bus, so a poll is the trigger. Each tick is one indexed
        SELECT and a text compare; the sites are only touched when the
        rendering changed.
        """
        await self._refresh_ip_whitelist_once()

    async def _refresh_ip_whitelist_once(self):
        # Serialised so the poll and the ConfigUpdate-triggered refresh
        # never write the same site twice at once.
        async with self._ip_whitelist_lock:
            try:
                await self._export_ip_whitelist()
            except OperationalError as exc:
                if "no such table" in str(exc):
                    # First tick after agent start can run before the
                    # resident schema is attached; the next tick catches up.
                    logger.info(
                        "IP whitelist export postponed: database not"
                        " ready yet (%s)",
                        exc,
                    )
                else:
                    logger.warning(
                        "IP whitelist export failed: %s", exc, exc_info=True
                    )
            except Exception as exc:
                # recurring_check abandons the loop after repeated errors;
                # a transient DB or disk problem must not switch the export
                # off until the agent restarts.
                logger.warning(
                    "IP whitelist export failed: %s", exc, exc_info=True
                )

    async def _export_ip_whitelist(self):
        if not Wordpress.SECURITY_PLUGIN_ENABLED:
            return
        if not self.installation_completed:
            return
        loaded = await plugin.load_ip_whitelist_php()
        if loaded is None:
            # AI Bot Protection is off (or no im360): nothing to export,
            # and an existing file is inert, so leave it alone.
            return
        text, count = loaded
        if text == self._last_ip_whitelist:
            logger.debug("IP whitelist unchanged, skipping export")
            return
        if count > 1000 and not self._ip_whitelist_size_warned:
            logger.warning(
                "IP whitelist has %d entries; exporting all of them to"
                " every WordPress site",
                count,
            )
        self._ip_whitelist_size_warned = count > 1000
        sites = await self._loop.run_in_executor(None, get_installed_sites)
        handled = removed = 0
        if sites:
            handled, removed = await plugin.sync_ip_whitelist_on_sites(
                sites, text, sink=self._sink
            )
        if text:
            logger.info(
                "Exported %d whitelist entries to %s on %d site(s)",
                count,
                plugin.IP_WHITELIST_FILENAME,
                handled,
            )
        elif removed:
            logger.info(
                "Removed %s from %d site(s) (whitelist empty)",
                plugin.IP_WHITELIST_FILENAME,
                removed,
            )
        else:
            # First empty tick after a restart advances the marker from None
            # to "" and walks every site, but the file usually never existed:
            # do not tell support the agent deleted 500 files it did not.
            logger.debug(
                "IP whitelist empty; %s already absent on all %d site(s)",
                plugin.IP_WHITELIST_FILENAME,
                handled,
            )
        # Dispatch marker: advanced after one pass regardless of per-site
        # outcome. A site that failed converges through the per-site data
        # writers (installer, scan hooks, daily update), as plugin_config
        # does; the poll never re-walks the fleet for one bad site.
        self._last_ip_whitelist = text

    @expect(MessageType.ConfigUpdate)
    async def manage_waf_config(self, message):
        """Caches are dispatch markers, not apply receipts — apply failures propagate, no auto-retry."""
        if not Wordpress.SECURITY_PLUGIN_ENABLED:
            return
        conf = message["conf"]
        config_dict = conf.config_to_dict()
        waf_value = config_dict.get("WORDPRESS", {}).get("waf_enabled")
        if isinstance(conf, UserConfig):
            if waf_value is None:
                prev = self._last_user_waf_enabled.pop(conf.username, None)
                if prev is None:
                    return
                new_effective = await plugin.is_waf_enabled_for_user(
                    conf.username
                )
                if not prev and new_effective:
                    await plugin.redeploy_waf_for_user(
                        conf.username, sink=self._sink
                    )
                elif prev and not new_effective:
                    await plugin.remove_waf_rules_for_user(conf.username)
                return
            if waf_value == self._last_user_waf_enabled.get(conf.username):
                return
            self._last_user_waf_enabled[conf.username] = waf_value
            if not waf_value or not plugin._get_global_waf_enabled():
                await plugin.remove_waf_rules_for_user(conf.username)
            else:
                await plugin.redeploy_waf_for_user(
                    conf.username, sink=self._sink
                )
        elif isinstance(conf, SystemConfig):
            try:
                current = Wordpress.WAF_ENABLED
            except KeyError:
                pass
            else:
                if current != self._last_waf_enabled:
                    self._last_waf_enabled = current
                    if not current:
                        await plugin.remove_waf_rules_for_all_sites()
                        save_state(WAF_SITES_STATE, {"sites": 0})
                    else:
                        await plugin.redeploy_waf_for_all_sites()

            try:
                current_default = Wordpress.WAF_DEFAULT
            except KeyError:
                pass
            else:
                if current_default != self._last_waf_default:
                    self._last_waf_default = current_default
                    await plugin.apply_waf_default_change(sink=self._sink)

    @expect(HookEvent.MalwareCleanupFinished)
    async def handle_malware_cleanup_finished(self, message):
        """
        INFO    [2025-02-24 12:00:20,384] imav.plugins.wordpress: Malware cleanup finished:
        HookEvent.MalwareCleanupFinished(
            {
                'cleanup_id': 'fa4fe7e48dbf45588f53b24366cd8893',
                'started': 1740398411.786418,
                'error': None,
                'total_files': 3,
                'total_cleaned': 3,
                'status': 'ok'
            }
        )
        """
        # Skip if plugin is disabled
        if not self.last_config_value:
            return

        # Leave early if status is not ok or the started time is missing.
        if message.get("status") != "ok" or not message.get("started"):
            return

        # load all malware hits cleaned since the cleanup started
        hits = _get_cleaned_malware_hits(message["started"])

        site_paths = set()

        # Collect all site paths that need to be updated.
        for hit in hits:
            if hit.resource_type == "file":
                try:
                    user_info = pwd.getpwnam(hit.user)
                    user_sites = get_sites_for_user(user_info)
                    uid = (  # In None cases there also no user_sites, so it wouldn't be used
                        user_info.pw_uid if user_info else None
                    )
                    for site_path in user_sites:
                        if hit.orig_file.startswith(site_path):
                            site_paths.add((site_path, uid))
                            break

                except KeyError:
                    pass

        if not site_paths:
            logger.debug("Cleanup finished => no sites found for cleaned hits")
            return

        logger.info(
            "Cleanup finished => %s site(s) need to be updated",
            len(site_paths),
        )

        # Convert paths to WPSite objects with empty domain and update data on the sites that need to be updated.
        # We need to work with paths here because sometimes the domain is not set, see https://cloudlinux.atlassian.net/browse/DEF-32238.
        wordpress_sites = [
            WPSite(docroot=site_path, domain="", uid=uid)
            for site_path, uid in site_paths
        ]

        await plugin.update_data_on_sites(self._sink, wordpress_sites)

        logger.info("%s site(s) updated after a cleanup", len(wordpress_sites))

    @expect(HookEvent.MalwareScanningFinished)
    async def handle_malware_scan_finished(self, message):
        """
        INFO    [2025-02-24 11:57:17,968] imav.plugins.wordpress: Malware scan finished:
        HookEvent.MalwareScanningFinished(
            {
                'scan_id': 'b9bd136aff0a4d87a248c859cfe41c47',
                'scan_type': 'user',
                'path': '/home/user1'
            }
        )
        INFO    [2025-02-24 12:00:10,740] imav.plugins.wordpress: Malware scan finished:
        HookEvent.MalwareScanningFinished(
            {
                'scan_id': 'a74271d2cdd04e0c9bd49ef6de23e0d8',
                'scan_type': 'user',
                'path': '/home/user4',
                'started': 1740398383,
                'total_files': 39229,
                'total_malicious': 3,
                'error': None,
                'status': 'ok',
                'scan_params': {'intensity_cpu': 2, 'intensity_io': 2, 'intensity_ram': 2048, 'initiator': None, 'file_patterns': None, 'exclude_patterns': None, 'follow_symlinks': False, 'detect_elf': True},
                'stats': {'scan_time': 27, 'mem_peak': 28217344, 'smart_time_hs': 0.004, 'scan_time_hs': 1.1751, 'smart_time_preg': 0, 'scan_time_preg': 2.7391, 'finder_time': 13.5896, 'cas_time': 0.7562, 'deobfuscate_time': 0.8998, 'total_files': 39229}
            }
        )
        """
        # Skip if plugin is disabled
        if not self.last_config_value:
            return

        # Leave early if status is not ok or path or stats are missing.
        if (
            message.get("status") != "ok"
            or not message.get("path")
            or not message.get("stats")
        ):
            return

        # Malware scan is finished, figure out what sites need to be updated based on the path.
        path = message["path"]
        sites = get_sites_by_path(path)
        if not sites:
            logger.debug("Scan finished => no sites found for path=%s", path)
            return

        # Update data on the sites that need to be updated.
        logger.info(
            "Scan finished => %s site(s) need to be updated", len(sites)
        )

        await plugin.update_data_on_sites(self._sink, sites)

        logger.info("%s site(s) updated after a scan", len(sites))
