import asyncio
import enum
import errno
import logging
import os
import pwd
import stat
import time

from collections import defaultdict
from collections.abc import Awaitable, Callable
from distutils.version import LooseVersion
from functools import cache
from pathlib import Path

from defence360agent.api import inactivity
from defence360agent.contracts.config import (
    MalwareScanScheduleInterval as Interval,
    SystemConfig,
    ANTIVIRUS_MODE,
    UserType,
    choose_value_from_config,
)
from defence360agent.files import Index, WP_RULES
from defence360agent.sentry import log_message
from defence360agent.utils import importer
from defence360agent.utils.fd_ops import (
    open_dir_no_symlinks,
    open_nofollow,
    rmtree_fd,
    safe_dir,
)
from defence360agent.contracts.config import Wordpress
from defence360agent.subsys.panels import hosting_panel
from defence360agent.wordpress.wp_rules import (
    get_wp_rules_data,
    get_wp_ruleset_version,
)
from defence360agent.model.wordpress import WordpressSite, WPSite
from defence360agent.model.wp_disabled_rule import WPDisabledRule
from defence360agent.wordpress import cli, telemetry
from defence360agent.wordpress.constants import (
    PLUGIN_SLUG,
    PLUGIN_VERSION_FILE,
)
from defence360agent.wordpress.utils import (
    _validate_preset,
    bucket_ip_whitelist_entries,
    calculate_next_scan_timestamp,
    clear_get_cagefs_enabled_users_cache,
    ensure_site_data_directory,
    fetch_manual_whitelist_entries,
    format_php_ip_whitelist,
    format_php_with_embedded_json,
    get_imunify_package_versions,
    get_last_scan,
    get_malware_history,
    prepare_plugin_config,
    prepare_scan_data,
    write_plugin_data_file_atomically,
)
from defence360agent.wordpress.site_repository import (
    clear_manually_deleted_flag,
    delete_site,
    get_installed_sites_by_domains,
    get_outdated_sites,
    get_sites_for_user,
    get_sites_to_adopt,
    get_sites_to_install,
    get_sites_to_mark_as_manually_deleted,
    get_installed_sites,
    insert_installed_sites,
    mark_site_as_manually_deleted,
    update_site_identity,
    update_site_version,
)

from defence360agent.wordpress.proxy_auth import setup_site_authentication

logger = logging.getLogger(__name__)


@cache
def _get_user_schedule_config_imav():
    return importer.get(
        module="imav.malwarelib.plugins.schedule_watcher",
        name="get_user_schedule_config",
        default=None,
    )


# Fallback when WORDPRESS keys are missing from config (old schema).
# True keeps WAF on for old schemas — no behavior change on upgrade.
_LEGACY_WAF_FALLBACK = True

_apply_waf_default_lock = asyncio.Lock()
_apply_waf_default_pending = False

_apply_ai_bot_default_lock = asyncio.Lock()
_apply_ai_bot_default_pending = False

# Default when ai_bot_protection is missing from config (old schema on
# a host where sibling packages haven't shipped the field yet). False
# because the feature is opt-in — if we can't determine admin intent,
# stay off rather than silently activate request-blocking logic.
_AI_BOT_PROTECTION_DEFAULT = False
_AI_BOT_PROTECTION_PRESET_DEFAULT = "balanced"


def _get_global_waf_enabled() -> bool:
    try:
        return bool(Wordpress.WAF_ENABLED)
    except KeyError:
        return _LEGACY_WAF_FALLBACK


def _get_waf_default() -> bool:
    try:
        return bool(Wordpress.WAF_DEFAULT)
    except KeyError:
        return _LEGACY_WAF_FALLBACK


def _get_security_plugin_enabled() -> bool:
    # KeyError -> False (feature off), matching the schema default. Unlike the
    # WAF keys there is no legacy "on" fallback: an absent key means the
    # feature simply isn't present on this schema, and a missing key must not
    # crash a read-only caller during agent/imunify-antivirus version skew.
    try:
        return bool(Wordpress.SECURITY_PLUGIN_ENABLED)
    except KeyError:
        return False


def _get_global_ai_bot_protection() -> bool:
    """Read WORDPRESS.ai_bot_protection from config, defaulting to False.

    Returns _AI_BOT_PROTECTION_DEFAULT when the config key is missing
    — e.g. the ai_bot_protection field hasn't rolled out to this
    install's imunify360 yet, or a sibling package is still on an
    older schema. Keeps the feature off in all ambiguous cases.
    """
    try:
        return bool(Wordpress.AI_BOT_PROTECTION)
    except KeyError:
        return _AI_BOT_PROTECTION_DEFAULT


def _get_global_ai_bot_protection_preset() -> str:
    """Read WORDPRESS.ai_bot_protection_preset from config, defaulting to
    "balanced".

    Two layers of safety: KeyError on a missing key (older schema, agent
    upgrade in progress) and _validate_preset() on the value itself
    (hand-edited override file, future preset rolled in via a sibling
    package this version doesn't recognise). Both fall back to the same
    canonical default so all layers — schema, agent, plugin — agree.
    """
    try:
        raw = Wordpress.AI_BOT_PROTECTION_PRESET
    except KeyError:
        return _AI_BOT_PROTECTION_PRESET_DEFAULT
    return _validate_preset(raw)


WAF_SOURCE_DEFAULT = "default"
WAF_SOURCE_OVERRIDE = "override"
WAF_SOURCE_KILL_SWITCH = "global kill switch"


def waf_global_snapshot() -> tuple[bool, bool, bool]:
    """Read the three server-wide WAF flags in one call.

    Returns (security_plugin_enabled, global_waf_enabled, waf_default), each
    guarded against a missing config key (schema version skew during an
    agent/imunify-antivirus upgrade) the same way the individual accessors are.
    """
    return (
        _get_security_plugin_enabled(),
        _get_global_waf_enabled(),
        _get_waf_default(),
    )


def waf_status_and_source_for_user_sync(username: str) -> tuple[bool, str]:
    if not _get_global_waf_enabled():
        return False, WAF_SOURCE_KILL_SWITCH
    try:
        value, source = choose_value_from_config(
            "WORDPRESS", "waf_enabled", username=username
        )
    except KeyError:
        return _get_waf_default(), WAF_SOURCE_DEFAULT
    if source == UserType.ROOT:
        return _get_waf_default(), WAF_SOURCE_DEFAULT
    return bool(value), WAF_SOURCE_OVERRIDE


def _is_waf_enabled_for_user_sync(username: str) -> bool:
    enabled, _ = waf_status_and_source_for_user_sync(username)
    return enabled


async def is_waf_enabled_for_user(username: str) -> bool:
    """Async wrapper — runs config file I/O in executor."""
    loop = asyncio.get_running_loop()
    return await loop.run_in_executor(
        None, _is_waf_enabled_for_user_sync, username
    )


async def _waf_enabled_or_default(username: str) -> bool:
    """WAF status for *username*; fail-open, defaults to enabled."""
    try:
        return await is_waf_enabled_for_user(username)
    except Exception:
        logger.warning(
            "Could not check WAF status for user %s, assuming enabled",
            username,
            exc_info=True,
        )
        return True


def _user_has_explicit_waf_override_sync(username: str) -> bool:
    try:
        _, source = choose_value_from_config(
            "WORDPRESS", "waf_enabled", username=username
        )
    except KeyError:
        return False
    return source != UserType.ROOT


def _get_ai_bot_protection_default() -> bool:
    # Unset (None) or key absent (imunify-antivirus predates it): fall back to
    # the global switch, the pre-split source of truth, so a no-override
    # account keeps its current state without a seed migration. An explicit
    # True/False is the admin's new-account default.
    try:
        value = Wordpress.AI_BOT_PROTECTION_DEFAULT
    except KeyError:
        return _get_global_ai_bot_protection()
    if value is None:
        return _get_global_ai_bot_protection()
    return bool(value)


AI_BOT_SOURCE_DEFAULT = "default"
AI_BOT_SOURCE_OVERRIDE = "override"
AI_BOT_SOURCE_KILL_SWITCH = "global kill switch"


def ai_bot_status_and_source_for_user_sync(username: str) -> tuple[bool, str]:
    if not _get_global_ai_bot_protection():
        return False, AI_BOT_SOURCE_KILL_SWITCH
    try:
        value, source = choose_value_from_config(
            "WORDPRESS", "ai_bot_protection", username=username
        )
    except KeyError:
        return _get_ai_bot_protection_default(), AI_BOT_SOURCE_DEFAULT
    if source == UserType.ROOT:
        return _get_ai_bot_protection_default(), AI_BOT_SOURCE_DEFAULT
    return bool(value), AI_BOT_SOURCE_OVERRIDE


def _is_ai_bot_protection_enabled_for_user_sync(username: str) -> bool:
    enabled, _ = ai_bot_status_and_source_for_user_sync(username)
    return enabled


def _user_has_explicit_ai_bot_override_sync(username: str) -> bool:
    try:
        _, source = choose_value_from_config(
            "WORDPRESS", "ai_bot_protection", username=username
        )
    except KeyError:
        return False
    return source != UserType.ROOT


def _no_override_uids_sync(usernames) -> set:
    """Uids of users without an explicit ai_bot_protection override.

    Sync on purpose: config reads and NSS lookups both block, so the whole
    walk runs in one executor call.
    """
    uids = set()
    for username in usernames:
        try:
            if _user_has_explicit_ai_bot_override_sync(username):
                continue
            uids.add(pwd.getpwnam(username).pw_uid)
        except KeyError:
            continue
        except Exception as e:
            logger.warning(
                "Failed to apply ai_bot_protection_default change"
                " for user %s: %s",
                username,
                e,
            )
    return uids


def ai_bot_preset_for_user_sync(username: str) -> str:
    try:
        value, source = choose_value_from_config(
            "WORDPRESS", "ai_bot_protection_preset", username=username
        )
    except KeyError:
        return _get_global_ai_bot_protection_preset()
    if source == UserType.ROOT:
        return _get_global_ai_bot_protection_preset()
    return _validate_preset(value)


def ai_bot_global_snapshot() -> tuple[bool, bool, bool]:
    """The three server-wide AI-Bot flags in one call: (security plugin
    enabled, global kill switch, effective new-account default). Each guarded
    against a missing config key the same way the individual accessors are."""
    return (
        _get_security_plugin_enabled(),
        _get_global_ai_bot_protection(),
        _get_ai_bot_protection_default(),
    )


COMPONENTS_DB_PATH = Path(
    "/var/lib/cloudlinux-app-version-detector/components_versions.sqlite3"
)


def _get_user_schedule_config(username: str, admin_config: SystemConfig):
    """
    Get user-specific schedule configuration with lazy import fallback.

    Returns default values if imav.malwarelib is not available.
    """
    get_user_schedule_config = _get_user_schedule_config_imav()
    if get_user_schedule_config is None:
        logger.debug(
            "imav.malwarelib not available, returning default schedule config"
        )
        return Interval.NONE, 0, 1, 0
    return get_user_schedule_config(username, admin_config)


def get_updated_wp_rules_data(index: Index) -> dict | None:
    """
    Retrieve WordPress rules with ANTIVIRUS_MODE handling and global disable filtering.

    In ANTIVIRUS_MODE, all rules are set to monitoring mode ("pass").
    Globally disabled rules are filtered out entirely — they should not
    appear in rules.php. Domain-specific disables are handled separately
    via disabled-rules.php.

    Args:
        index: The Index object used to locate the wp-rules.zip file.

    Returns:
        The parsed wp-rules data with mode adjusted for ANTIVIRUS_MODE
        and globally disabled rules removed, or None if rules cannot be loaded.
    """
    rules_data = get_wp_rules_data(index)
    if rules_data is None:
        return None

    if ANTIVIRUS_MODE:
        # all rules will be in monitoring mode only for AV and AV+
        for cve, params in rules_data.items():
            params["mode"] = "pass"

    # Filter out globally disabled rules — these are excluded from rules.php
    globally_disabled = set(WPDisabledRule.get_global_disabled())
    if globally_disabled:
        rules_data = {
            cve: params
            for cve, params in rules_data.items()
            if cve not in globally_disabled
        }

    return rules_data


def clear_caches():
    """Clear all WordPress-related caches."""
    clear_get_cagefs_enabled_users_cache()
    cli.clear_get_content_dir_cache()


def site_search(items: dict, user_info: pwd.struct_passwd, matcher) -> dict:
    # Get all WordPress sites for the user (the main site is always last)
    user_sites = get_sites_for_user(user_info)
    result = {path: [] for path in user_sites}
    for item in items:
        # Find all matching sites for this item
        matching_sites = [path for path in user_sites if matcher(item, path)]

        if matching_sites:
            # Find the most specific (longest) matching path
            most_specific_site = max(matching_sites, key=len)
            result[most_specific_site].append(item)

    return result


async def _get_scan_data_for_user(
    sink, user_info: pwd.struct_passwd, admin_config: SystemConfig
):
    # Get the last scan data
    last_scan = await get_last_scan(sink, user_info.pw_name)

    # Extract the last scan date
    last_scan_time = last_scan.get("scan_date", None)

    # Get user-specific schedule configuration
    interval, hour, day_of_month, day_of_week = _get_user_schedule_config(
        user_info.pw_name, admin_config
    )

    next_scan_time = None
    if interval != Interval.NONE:
        next_scan_time = calculate_next_scan_timestamp(
            interval, hour, day_of_month, day_of_week
        )

    # Get the malware history for the user
    malware_history = get_malware_history(user_info.pw_name)

    # Split malware history by site. This part relies on the main site being the last one in the list.
    # Without this all malware could be attributed to the main site.
    malware_by_site = site_search(
        malware_history,
        user_info,
        lambda item, path: item["resource_type"] == "file"
        and item["file"].startswith(path),
    )

    return last_scan_time, next_scan_time, malware_by_site


async def _send_telemetry_task(coro, semaphore: asyncio.Semaphore):
    async with semaphore:
        try:
            await coro
        except Exception as e:
            logger.error(f"Telemetry task failed: {e}")


async def process_telemetry_tasks(coroutines: list, concurrency=10):
    """
    Process a list of telemetry coroutines with a concurrency limit.s
    """
    if not coroutines:
        return

    semaphore = asyncio.Semaphore(concurrency)
    tasks = [
        asyncio.create_task(_send_telemetry_task(coro, semaphore))
        for coro in coroutines
    ]

    try:
        await asyncio.gather(*tasks)
    except Exception as e:
        logger.error(f"Some telemetry tasks failed: {e}")


async def load_wp_rules_php():
    """
    Load WordPress rules from the index and format them as PHP.

    Returns:
        str or None: PHP-formatted rules data, or None if rules could not be loaded.
    """
    try:
        wp_rules_index = Index(WP_RULES, integrity_check=False)
        await wp_rules_index.update()
        wp_rules_data = get_updated_wp_rules_data(wp_rules_index)
    except Exception as e:
        logger.warning(
            "Failed to load wp-rules index: %s, skipping rules installation",
            e,
        )
        return None

    if not wp_rules_data:
        logger.warning(
            "valid WordPress rules not found, skipping rules installation"
        )
        return None

    # Get version and create ruleset dict with version and rules
    wp_rules_version = get_wp_ruleset_version(wp_rules_index)
    ruleset_dict = {
        "version": wp_rules_version,
        "rules": wp_rules_data,
    }
    return format_php_with_embedded_json(ruleset_dict)


async def install_everywhere(sink):
    """Install the imunify-security plugin for all sites where it is not installed."""
    sites = get_sites_to_install()
    installer = WordPressSiteInstaller(sink, sites)
    return await installer.run()


async def install_on_site(sink, site: WPSite):
    """Install the imunify-security plugin on exactly one site.

    Used by the restore-site CLI: the fleet-wide install_everywhere pass is
    considerably heavier than needed for a single docroot.
    """
    installer = WordPressSiteInstaller(sink, {site})
    return await installer.run()


async def adopt_found_sites(sink):
    """
    Adopt WordPress sites where the plugin is installed but not tracked in our database
    or flagged as manually removed.

    This handles scenarios like:
    - Sites copied/migrated from another location
    - Sites migrated from another server
    - Sites where the manually_deleted flag was incorrectly set (past bugs)
    - Sites where the user installed the plugin from wordpress.org
    """
    sites = get_sites_to_adopt()
    processor = WordPressSiteAdopter(sink, sites)
    return await processor.run()


def get_latest_plugin_version() -> str:
    """Get the latest version of the imunify-security plugin from the version file."""
    try:
        if not PLUGIN_VERSION_FILE.exists():
            logger.error(
                "Plugin version file does not exist: %s", PLUGIN_VERSION_FILE
            )
            return None
        return PLUGIN_VERSION_FILE.read_text().strip()
    except Exception as e:
        logger.error("Failed to read plugin version file: %s", e)
        return None


async def update_everywhere(sink):
    """Update the imunify-security plugin on all sites where it is installed."""
    latest_version = get_latest_plugin_version()
    if not latest_version:
        logger.error("Could not determine latest plugin version")
        return

    logger.info(
        "Updating imunify-security wp plugin to the latest version %s",
        latest_version,
    )

    updated = set()
    telemetry_coros = []
    with inactivity.track.task("wp-plugin-update"):
        try:
            # Get sites with outdated versions
            outdated_sites = get_outdated_sites(latest_version)
            logger.info(f"Found {len(outdated_sites)} outdated sites")

            if not outdated_sites:
                return

            # Create SystemConfig once for all users
            admin_config = SystemConfig()

            versions = await get_imunify_package_versions()

            # Group sites by user id
            sites_by_user = defaultdict(list)
            for site in outdated_sites:
                sites_by_user[site.uid].append(site)

            # Process each user's sites
            for uid, sites in sites_by_user.items():
                try:
                    user_info = pwd.getpwuid(uid)
                    username = user_info.pw_name
                except Exception as error:
                    logger.error(
                        "Failed to get username for uid=%d. error=%s",
                        uid,
                        error,
                    )
                    continue

                # Get scan data once for all sites of this user
                (
                    last_scan_time,
                    next_scan_time,
                    malware_by_site,
                ) = await _get_scan_data_for_user(
                    sink, user_info, admin_config
                )
                plugin_config = prepare_plugin_config(username)

                for site in sites:
                    if await remove_site_if_missing(sink, site):
                        continue
                    try:
                        # Check if site still exists
                        if not await cli.is_wordpress_installed(site):
                            logger.info(
                                "WordPress site no longer exists: %s", site
                            )
                            continue

                        # Prepare scan data
                        scan_data = prepare_scan_data(
                            last_scan_time,
                            next_scan_time,
                            username,
                            site,
                            malware_by_site,
                            versions=versions,
                        )

                        # Resolve the data dir once; both writes reuse it.
                        data_dir = await ensure_site_data_directory(
                            site, user_info
                        )
                        if data_dir is None:
                            continue

                        # Update the scan data file
                        await update_scan_data_file(
                            site,
                            scan_data,
                            user_info=user_info,
                            data_dir=data_dir,
                        )

                        # Keep plugin_config.php fresh alongside scan_data
                        # — covers the case where a ConfigUpdate event
                        # was missed (plugin re-installed after the
                        # toggle, first scan after an upgrade, etc).
                        await update_plugin_config_file(
                            site,
                            plugin_config,
                            user_info=user_info,
                            data_dir=data_dir,
                        )
                        ip_whitelist = await _load_ip_whitelist_or_skip()
                        if ip_whitelist is not None:
                            text, _ = ip_whitelist
                            await _sync_ip_whitelist_or_skip(
                                site,
                                text,
                                user_info=user_info,
                                data_dir=data_dir,
                            )

                        # Now update the plugin
                        await cli.plugin_update(site)
                        updated.add(site)

                        # Get the version after update
                        version = await cli.get_plugin_version(site)
                        if version:
                            # Store original version for comparison
                            original_version = site.version

                            # Update the database with the new version
                            update_site_version(site, version)

                            # Create a new WPSite with updated version
                            site = site.build_with_version(version)

                            # Determine if this is a downgrade
                            is_downgrade = LooseVersion(
                                version
                            ) < LooseVersion(original_version)

                            # Prepare telemetry
                            telemetry_coros.append(
                                telemetry.send_event(
                                    sink=sink,
                                    event=(
                                        "downgraded_by_imunify"
                                        if is_downgrade
                                        else "updated_by_imunify"
                                    ),
                                    site=site,
                                    version=version,
                                )
                            )

                    except Exception as error:
                        logger.error(
                            "Failed to update plugin on site=%s error=%s",
                            site,
                            error,
                        )

            logger.info(
                "Updated imunify-security wp plugin on %d sites",
                len(updated),
            )
        except asyncio.CancelledError:
            logger.info(
                "Update of imunify-security wp plugin was cancelled. Plugin"
                " was updated on %d sites",
                len(updated),
            )
        except Exception as error:
            logger.error(
                "Error occurred during plugin update. error=%s", error
            )
            raise
        finally:
            # Send telemetry
            await process_telemetry_tasks(telemetry_coros)


async def delete_plugin_files(site: WPSite):
    data_dir = await cli.get_data_dir(site)
    # Open both target and parent dirs with symlink protection before
    # performing any destructive operations.
    try:
        dir_fd = open_dir_no_symlinks(data_dir)
    except OSError as exc:
        if exc.errno == errno.ENOENT:
            return  # directory does not exist — nothing to delete
        if exc.errno in (errno.ELOOP, errno.ENOTDIR):
            logger.warning(
                "Skipping rmtree: data directory %s is a symlink", data_dir
            )
            return
        raise

    try:
        with safe_dir(data_dir.parent) as parent_fd:
            try:
                await asyncio.to_thread(rmtree_fd, dir_fd)
            finally:
                os.close(dir_fd)
                dir_fd = -1
            # Remove the now-empty directory via the parent fd.
            os.rmdir(data_dir.name, dir_fd=parent_fd)
    except BaseException:
        if dir_fd >= 0:
            os.close(dir_fd)
        raise


async def remove_from_single_site(site: WPSite, sink, telemetry_coros) -> int:
    """
    Remove the imunify-security plugin from a single site, including all cleanup and telemetry.
    Returns the number of affected sites (should be 1 if deletion was successful).
    This function is intended to be protected with asyncio.shield to ensure it completes even if the parent task is cancelled.
    """
    try:
        # Check if site is still installed and accessible using WP CLI
        is_installed = await cli.is_plugin_installed(site)
        if not is_installed:
            # Plugin is no longer installed. It was removed manually by the user.
            await process_manually_deleted_plugin(
                site, time.time(), sink, telemetry_coros
            )
            return 0

        # Get the version of the plugin (for telemetry data)
        version = await cli.get_plugin_version(site)

        # Uninstall the plugin from WordPress site.
        await cli.plugin_uninstall(site)

        # Delete the data files from the site.
        await delete_plugin_files(site)

        # Delete the site from database.
        affected = delete_site(site)

        # Send telemetry for successful uninstall
        telemetry_coros.append(
            telemetry.send_event(
                sink=sink,
                event="uninstalled_by_imunify",
                site=site,
                version=version,
            )
        )
        return affected
    except Exception as error:
        # Log any error that occurs during the removal process
        logger.error("Failed to remove plugin from %s %s", site, error)
        return 0


async def remove_all_installed(sink):
    """Remove the imunify-security plugin from all sites where it is installed."""
    logger.info("Deleting imunify-security wp plugin")

    telemetry_coros = []
    affected = 0
    with inactivity.track.task("wp-plugin-removal"):
        try:
            clear_caches()

            to_remove = get_installed_sites()

            for site in to_remove:
                try:
                    affected += await asyncio.shield(
                        remove_from_single_site(site, sink, telemetry_coros)
                    )
                except asyncio.CancelledError:
                    logger.info(
                        "Deleting imunify-security wp plugin was cancelled."
                        " Plugin was deleted from %d sites (out of %d)",
                        affected,
                        len(to_remove),
                    )
        except Exception as error:
            logger.error("Error occurred during plugin deleting. %s", error)
            raise
        finally:
            logger.info(
                "Removed imunify-security wp plugin from %s sites",
                affected,
            )

            #  send telemetry
            await process_telemetry_tasks(telemetry_coros)


async def _is_plugin_present_on_disk(site: WPSite) -> bool:
    # os.path.isdir() is not used here: it reports every stat failure as
    # "not a directory", so an unreadable docroot would read as a deleted
    # site and take the tidy-up branch before any later check runs.
    try:
        docroot_stat = os.stat(site.docroot)
    except (FileNotFoundError, NotADirectoryError):
        return False
    except (OSError, ValueError) as error:
        # ValueError too: os.stat rejects an embedded NUL, and both the
        # isdir()/exists() calls this replaced reported that as "absent".
        logger.warning(
            "Failed to stat %s, assuming the plugin is still installed: %s",
            site.docroot,
            error,
        )
        return True

    if not stat.S_ISDIR(docroot_stat.st_mode):
        return False

    try:
        content_dir = await cli.get_content_dir(site)
    except Exception as error:
        # Uncertainty must not take the delete branch: data may be in use.
        logger.warning(
            "Failed to resolve the content directory of %s, assuming the"
            " plugin is still installed: %s",
            site.docroot,
            error,
        )
        return True

    if not os.path.isdir(content_dir):
        # get_content_dir() falls back to the default path when WP CLI fails,
        # so a missing directory means unknown location, not a gone plugin.
        logger.warning(
            "Content directory %s of %s does not exist, assuming the plugin is"
            " still installed",
            content_dir,
            site.docroot,
        )
        return True

    plugin_file = content_dir / "plugins" / PLUGIN_SLUG / f"{PLUGIN_SLUG}.php"
    # Path.exists() is not used here: it hides ELOOP, so an unresolvable
    # symlink would read as a gone plugin, and it raises every other error
    # into the caller, which gives up on the rest of the batch.
    try:
        os.stat(plugin_file)
    except (FileNotFoundError, NotADirectoryError):
        return False
    except (OSError, ValueError) as error:
        logger.warning(
            "Failed to stat %s, assuming the plugin is still installed: %s",
            plugin_file,
            error,
        )
        return True

    return True


async def process_manually_deleted_plugin(site, now, sink, telemetry_coros):
    """
    Process the manually deleted plugin for a single site.

    Args:
        site: The site to process.
        now: The current time.
        sink: The telemetry/event sink.
        telemetry_coros: The list of telemetry coroutines to add the event to.

    The process includes:
    - marking the site as manually deleted in the database
    - removing plugin data files
    - sending telemetry for manual removal
    """
    try:
        # Mark the site as manually deleted in the database
        mark_site_as_manually_deleted(site, now)

        # Remove plugin data files
        await delete_plugin_files(site)

        # Send telemetry for manual removal
        telemetry_coros.append(
            telemetry.send_event(
                sink=sink,
                event="removed_by_user",
                site=site,
                version=site.version,
            )
        )
    except Exception as error:
        logger.error(
            "Failed to process manually deleted plugin for site=%s error=%s",
            site,
            error,
        )


async def tidy_up_manually_deleted(
    sink, freshly_installed_sites: set[WPSite] = None
):
    """
    Tidy up sites that have been manually deleted by the user.

    Args:
        sink: The telemetry/event sink.
        freshly_installed_sites: Optional set of sites that were just installed and should be excluded
                                from being marked as manually deleted to avoid race conditions.
    """
    telemetry_coros = []
    try:
        to_mark_as_manually_removed = get_sites_to_mark_as_manually_deleted(
            freshly_installed_sites
        )
        if to_mark_as_manually_removed:
            now = time.time()
            for site in to_mark_as_manually_removed:
                # Candidates come from an AppVersionDetector snapshot that can
                # be up to a day older than the plugin installation.
                if await _is_plugin_present_on_disk(site):
                    logger.warning(
                        "Skipping tidy up of %s: the plugin is still present"
                        " on disk while the app detector reports it missing",
                        site.docroot,
                    )
                    continue

                await process_manually_deleted_plugin(
                    site, now, sink, telemetry_coros
                )

    except Exception as error:
        logger.error("Error occurred during site tidy up. %s", error)
    finally:
        if telemetry_coros:
            await process_telemetry_tasks(telemetry_coros)


async def update_data_on_sites(sink, sites: list[WPSite]):
    if not sites:
        return

    # Create SystemConfig once for all users
    admin_config = SystemConfig()

    versions = await get_imunify_package_versions()

    # Group sites by user id
    sites_by_user = defaultdict(list)
    for site in sites:
        sites_by_user[site.uid].append(site)

    # Now iterate over the grouped sites
    for uid, sites in sites_by_user.items():
        try:
            user_info = pwd.getpwuid(uid)
            username = user_info.pw_name
        except Exception as error:
            logger.error(
                "Failed to get username for uid=%d. error=%s",
                uid,
                error,
            )
            continue

        (
            last_scan_time,
            next_scan_time,
            malware_by_site,
        ) = await _get_scan_data_for_user(sink, user_info, admin_config)
        plugin_config = prepare_plugin_config(username)

        for site in sites:
            if await remove_site_if_missing(sink, site):
                continue
            try:
                # Prepare scan data
                scan_data = prepare_scan_data(
                    last_scan_time,
                    next_scan_time,
                    username,
                    site,
                    malware_by_site,
                    versions=versions,
                )

                # Resolve the site's data directory once; both writes reuse it.
                data_dir = await ensure_site_data_directory(site, user_info)
                if data_dir is None:
                    continue

                # Update the scan data file
                await update_scan_data_file(
                    site, scan_data, user_info=user_info, data_dir=data_dir
                )

                # Keep plugin_config.php fresh alongside scan_data.
                await update_plugin_config_file(
                    site, plugin_config, user_info=user_info, data_dir=data_dir
                )

                # Read right before the write: a pass that shells out for
                # minutes must not put back an entry deleted (and removed
                # by the poll) in the meantime. None: gated off.
                ip_whitelist = await _load_ip_whitelist_or_skip()
                if ip_whitelist is not None:
                    text, _ = ip_whitelist
                    await _sync_ip_whitelist_or_skip(
                        site,
                        text,
                        user_info=user_info,
                        data_dir=data_dir,
                    )
            except Exception as error:
                logger.error(
                    "Failed to update site data on site=%s error=%s",
                    site,
                    error,
                )


async def _write_json_php_data_file(
    site: WPSite,
    filename: str,
    data: dict,
    *,
    user_info: pwd.struct_passwd | None = None,
    data_dir: Path | None = None,
) -> bool:
    """Write ``data`` as embedded JSON to ``<site data dir>/<filename>``.

    A caller writing several files into one site's directory can resolve
    ``user_info`` and ``data_dir`` once and pass them in, so the owner
    lookup and directory-ensure are not repeated per file.

    Returns False when the site's data directory is unavailable and
    nothing was written, True after a successful write.
    """
    if user_info is None:
        user_info = pwd.getpwuid(site.uid)
    if data_dir is None:
        data_dir = await ensure_site_data_directory(site, user_info)
        if data_dir is None:
            return False
    php_content = format_php_with_embedded_json(data)
    write_plugin_data_file_atomically(
        data_dir / filename, php_content, uid=site.uid, gid=user_info.pw_gid
    )
    return True


async def update_scan_data_file(
    site: WPSite,
    scan_data: dict,
    *,
    user_info: pwd.struct_passwd | None = None,
    data_dir: Path | None = None,
):
    await _write_json_php_data_file(
        site,
        "scan_data.php",
        scan_data,
        user_info=user_info,
        data_dir=data_dir,
    )


async def update_plugin_config_file(
    site: WPSite,
    plugin_config: dict,
    *,
    user_info: pwd.struct_passwd | None = None,
    data_dir: Path | None = None,
) -> bool:
    """
    Write plugin_config.php for a single WordPress site.

    Separate file from scan_data.php so a config toggle doesn't force
    rewriting the malware list, and so the mu-plugin hot path loads
    only what it needs per request.

    Returns False when the write was skipped because the site's data
    directory is unavailable, True otherwise.
    """
    return await _write_json_php_data_file(
        site,
        "plugin_config.php",
        plugin_config,
        user_info=user_info,
        data_dir=data_dir,
    )


IP_WHITELIST_FILENAME = "ip_whitelist.php"


async def load_ip_whitelist_php() -> tuple[str, int] | None:
    """Gate, fetch and render the ip_whitelist.php export in one step.

    Returns None when nothing should be touched: the global AI Bot
    Protection toggle is off (the plugin only reads the file inside the
    bot pipeline, so an existing export is inert and is left in place), or
    there is no im360 IPList model (AV-only install). Returns ("", 0) when
    the manual whitelist is empty, meaning an existing file must be
    removed, and otherwise (php_text, active_entry_count).
    """
    if not _get_global_ai_bot_protection():
        return None
    # On the loop thread on purpose: the resident schema is attached to the
    # main connection only (see fetch_manual_whitelist_entries).
    entries = fetch_manual_whitelist_entries()
    if entries is None:
        return None
    by_octet, broad = bucket_ip_whitelist_entries(entries, now=time.time())
    count = sum(len(group) for group in by_octet.values()) + len(broad)
    if count == 0:
        return "", 0
    return format_php_ip_whitelist(by_octet, broad), count


async def _load_ip_whitelist_or_skip() -> tuple[str, int] | None:
    """load_ip_whitelist_php for the per-site data-file writers.

    Those passes did not depend on the resident database before; a failure
    there (schema not attached yet at startup, a locked database) must cost
    only that site's whitelist export, not the scan_data/plugin_config
    writes or the whole installation.
    """
    try:
        return await load_ip_whitelist_php()
    except Exception as exc:
        logger.warning("IP whitelist export skipped for this pass: %s", exc)
        return None


async def _sync_ip_whitelist_or_skip(
    site: WPSite,
    php_text: str,
    *,
    user_info: pwd.struct_passwd,
    data_dir: Path,
) -> None:
    """sync_ip_whitelist_file for the per-site data-file writers.

    The export is the only write in those passes that always creates a new
    file — the others no-op when the bytes match — so it is the first to
    fail on a full filesystem or an account over quota. That must cost the
    site its whitelist, not its plugin install or update.
    """
    try:
        await sync_ip_whitelist_file(
            site, php_text, user_info=user_info, data_dir=data_dir
        )
    except Exception as exc:
        logger.warning(
            "IP whitelist export skipped for site=%s: %s", site, exc
        )


class SiteSyncResult(enum.Enum):
    """Outcome of one site's ip_whitelist.php write or removal.

    ``REMOVED`` means a file was actually unlinked; ``ABSENT`` means the site
    was already converged with nothing to do. Both count as handled, but only
    ``REMOVED`` is a real change worth reporting at INFO — the removal pass
    walks every site, most of which never had the file.
    """

    WROTE = "wrote"
    REMOVED = "removed"
    ABSENT = "absent"
    FAILED = "failed"


def _remove_ip_whitelist_file(
    data_dir: Path, docroot: str, uid: int
) -> SiteSyncResult:
    """Remove ip_whitelist.php from data_dir via a symlink-safe dir fd.

    Returns ``REMOVED`` when a file was unlinked, ``ABSENT`` when it was
    already gone (both are the target state), and ``FAILED`` when a guard or
    error left it in place.
    """
    try:
        dir_fd = open_dir_no_symlinks(data_dir)
    except FileNotFoundError:
        return SiteSyncResult.ABSENT
    except OSError as exc:
        if exc.errno in (errno.ELOOP, errno.ENOTDIR):
            logger.warning(
                "Skipping %s removal: data dir %s is a symlink or not a"
                " directory",
                IP_WHITELIST_FILENAME,
                data_dir,
            )
            return SiteSyncResult.FAILED
        logger.error("Failed to open data dir for %s: %s", docroot, exc)
        return SiteSyncResult.FAILED
    try:
        dir_uid = os.fstat(dir_fd).st_uid
        if dir_uid != uid:
            logger.warning(
                "Skipping %s removal: data dir %s for %s is owned by uid %s,"
                " expected %s",
                IP_WHITELIST_FILENAME,
                data_dir,
                docroot,
                dir_uid,
                uid,
            )
            return SiteSyncResult.FAILED
        try:
            os.remove(IP_WHITELIST_FILENAME, dir_fd=dir_fd)
        except FileNotFoundError:
            return SiteSyncResult.ABSENT
        except OSError as exc:
            logger.error(
                "Failed to remove %s from %s: %s",
                IP_WHITELIST_FILENAME,
                docroot,
                exc,
            )
            return SiteSyncResult.FAILED
        logger.debug("Removed %s from %s", IP_WHITELIST_FILENAME, docroot)
        return SiteSyncResult.REMOVED
    finally:
        os.close(dir_fd)


async def sync_ip_whitelist_file(
    site: WPSite,
    php_text: str,
    *,
    user_info: pwd.struct_passwd,
    data_dir: Path,
) -> SiteSyncResult:
    """Bring one site's ip_whitelist.php in line with *php_text*.

    Empty text means "no whitelist": the file is removed if present.
    Otherwise it is written atomically with the site owner and the same
    mode as plugin_config.php; an unchanged file is not rewritten.
    """
    if not php_text:
        return await asyncio.to_thread(
            _remove_ip_whitelist_file, data_dir, site.docroot, site.uid
        )
    await asyncio.to_thread(
        write_plugin_data_file_atomically,
        data_dir / IP_WHITELIST_FILENAME,
        php_text,
        uid=site.uid,
        gid=user_info.pw_gid,
    )
    return SiteSyncResult.WROTE


async def sync_ip_whitelist_on_sites(
    sites: list[WPSite], php_text: str, sink=None
) -> tuple[int, int]:
    """Write (or remove, for empty *php_text*) ip_whitelist.php on every site.

    Returns ``(handled, removed)``: *handled* is the number of sites brought
    to the target state — written, reaped, or skipped because their wp-content
    is unusable — reported the way update_plugin_config_on_sites does;
    *removed* is the subset where a file was actually unlinked, so an
    empty-whitelist pass over sites that never had the file reports zero real
    removals. The dispatch marker advances either way: a site that failed
    converges through the per-site data-file writers. Per-site failures are
    logged and not counted.
    """
    if not sites:
        return 0, 0

    handled = 0
    removed = 0
    sites_by_user: dict[int, list[WPSite]] = defaultdict(list)
    for site in sites:
        sites_by_user[site.uid].append(site)

    for uid, user_sites in sites_by_user.items():
        try:
            user_info = pwd.getpwuid(uid)
        except Exception as error:
            logger.error(
                "Failed to get username for uid=%d. error=%s", uid, error
            )
            continue

        for site in user_sites:
            try:
                if await remove_site_if_missing(sink, site):
                    handled += 1
                    continue
                if php_text:
                    data_dir = await ensure_site_data_directory(
                        site, user_info
                    )
                    if data_dir is None:
                        handled += 1
                        logger.warning(
                            "%s not synced for site=%s: its wp-content is"
                            " unusable; the site is skipped until"
                            " reinstalled",
                            IP_WHITELIST_FILENAME,
                            site,
                        )
                        continue
                else:
                    # Removal only unlinks. ensure_site_data_directory would
                    # mkdir and chown the directory and rewrite the listing
                    # protection files on every site of every server without
                    # manual entries, on every agent start.
                    data_dir = await cli.get_data_dir(site)
                result = await sync_ip_whitelist_file(
                    site, php_text, user_info=user_info, data_dir=data_dir
                )
                if result is not SiteSyncResult.FAILED:
                    handled += 1
                if result is SiteSyncResult.REMOVED:
                    removed += 1
            except Exception as error:
                logger.error(
                    "Failed to sync %s on site=%s error=%s",
                    IP_WHITELIST_FILENAME,
                    site,
                    error,
                )

    return handled, removed


async def update_plugin_config_on_sites(sites: list[WPSite], sink=None) -> int:
    """
    Rewrite plugin_config.php on every managed site in one pass.

    Used by the ConfigUpdate handler that reacts to
    WORDPRESS.ai_bot_protection toggles. Writes only plugin_config.php
    — scan_data.php is untouched, so a toggle doesn't churn the
    (potentially large) malware payload or wait on a scan cycle.

    A site whose docroot is gone is reaped (emitting site_removed on
    the sink, if given) before its write.

    Returns the number of sites handled — written, reaped, or skipped
    with a warning because wp-content is unusable — so the caller can
    decide whether to advance its cached state. Skipped sites count as
    handled deliberately: they cannot take a write until reinstalled
    (the installer writes a fresh plugin_config.php then), so leaving
    them out would keep the caller retrying the whole pass forever.
    """
    if not sites:
        return 0

    updated = 0
    skipped = 0

    # Group by uid so we look up username once per user, mirroring
    # update_data_on_sites' pattern and making per-user error isolation
    # straightforward.
    sites_by_user: dict[int, list[WPSite]] = defaultdict(list)
    for site in sites:
        sites_by_user[site.uid].append(site)

    for uid, user_sites in sites_by_user.items():
        try:
            user_info = pwd.getpwuid(uid)
            plugin_config = prepare_plugin_config(user_info.pw_name)
        except Exception as error:
            logger.error(
                "Failed to prepare plugin config for uid=%d. error=%s",
                uid,
                error,
            )
            continue

        for site in user_sites:
            try:
                if await remove_site_if_missing(sink, site):
                    updated += 1
                    continue
                if await update_plugin_config_file(
                    site, plugin_config, user_info=user_info
                ):
                    updated += 1
                else:
                    skipped += 1
                    logger.warning(
                        "plugin_config.php not written for site=%s: its"
                        " wp-content is unusable; the site keeps its old"
                        " config until it is reinstalled",
                        site,
                    )
            except Exception as error:
                logger.error(
                    "Failed to update plugin_config.php on site=%s error=%s",
                    site,
                    error,
                )

    return updated + skipped


async def update_wp_rules_for_site(
    site: WPSite,
    user_info: pwd.struct_passwd,
    wp_rules_php: str,
    updated: set,
    failed: set,
) -> None:
    """
    Deploy wp-rules to a single WordPress site and track the result.

    Args:
        site: WordPress site to deploy to
        user_info: User information from pwd
        wp_rules_php: Formatted PHP rules content
        updated: Set to add site to if successful
        failed: Set to add site to if failed
    """
    gid = user_info.pw_gid

    try:
        data_dir = await ensure_site_data_directory(site, user_info)
        if data_dir is None:
            return
        rules_path = data_dir / "rules.php"
        write_plugin_data_file_atomically(
            rules_path, wp_rules_php, uid=site.uid, gid=gid
        )
        updated.add(site)
        logger.info("Updated wp-rules for site %s", site.docroot)
    except Exception as error:
        failed.add(site)
        logger.error(
            "Failed to update wp-rules for site %s: %s",
            site.docroot,
            error,
        )


async def _deploy_to_sites(
    sites: list[WPSite],
    make_task: Callable[
        [WPSite, pwd.struct_passwd, set, set], Awaitable[None]
    ],
    task_name: str,
    fingerprint: str,
    sink=None,
) -> None:
    """
    Run a per-site async deployment over a list of WordPress sites.

    Groups sites by user, resolves UIDs, then runs tasks concurrently
    in batches.

    Args:
        sites: WordPress sites to deploy to
        make_task: Callable that creates a coroutine for one site.
            Signature: (site, user_info, updated_set, failed_set) -> awaitable
        task_name: Human-readable name for logging and inactivity tracking
        fingerprint: Sentry fingerprint for user-lookup failures
        sink: Optional telemetry sink for remove_site_if_missing
    """
    updated = set()
    failed = set()

    with inactivity.track.task(task_name):
        try:
            start_time = time.time()

            sites_by_user = defaultdict(list)
            for site in sites:
                sites_by_user[site.uid].append(site)

            tasks = []
            for uid, user_sites in sites_by_user.items():
                try:
                    user_info = pwd.getpwuid(uid)
                    username = user_info.pw_name
                except Exception as error:
                    log_message(
                        "Skipping {task} update for {count} site(s)"
                        " belonging to user {user} because username"
                        " retrieval failed. Reason: {reason}",
                        format_args={
                            "task": task_name,
                            "count": len(user_sites),
                            "user": uid,
                            "reason": error,
                        },
                        level="warning",
                        component="wordpress",
                        fingerprint=fingerprint,
                    )
                    for site in user_sites:
                        failed.add(site)
                    continue

                if not await _waf_enabled_or_default(username):
                    logger.info(
                        "WAF disabled for user %s, removing WAF files from"
                        " %d site(s)",
                        username,
                        len(user_sites),
                    )
                    await _remove_waf_files_from_sites(user_sites, username)
                    continue

                for site in user_sites:
                    if await remove_site_if_missing(sink, site):
                        continue
                    tasks.append(make_task(site, user_info, updated, failed))

            max_concurrent = 10
            for i in range(0, len(tasks), max_concurrent):
                batch = tasks[i : i + max_concurrent]
                await asyncio.gather(*batch, return_exceptions=True)

            elapsed = time.time() - start_time
            logger.info(
                "%s deployment complete. Updated: %d, Failed: %d,"
                " Duration: %.2fs",
                task_name,
                len(updated),
                len(failed),
                elapsed,
            )

        except asyncio.CancelledError:
            logger.info(
                "%s deployment was cancelled. Updated %d sites",
                task_name,
                len(updated),
            )
        except Exception as error:
            logger.error(
                "Error occurred during %s deployment. error=%s",
                task_name,
                error,
            )
            raise


async def _deploy_wp_rules_php(wp_rules_php: str, sink=None) -> None:
    """Deploy pre-formatted wp-rules PHP content to all active WordPress sites."""
    clear_caches()

    installed_sites = get_installed_sites()
    if not installed_sites:
        logger.debug("No active WordPress sites found")
        return

    def make_task(site, user_info, updated, failed):
        return update_wp_rules_for_site(
            site, user_info, wp_rules_php, updated, failed
        )

    await _deploy_to_sites(
        installed_sites,
        make_task,
        task_name="wp-rules",
        fingerprint="wp-rules-update-skip-user",
        sink=sink,
    )


async def update_wp_rules_on_sites(index: Index, is_updated: bool) -> None:
    """
    Hook that runs when wp-rules files are updated.
    Extracts wp-rules.yaml from wp-rules.zip and deploys to all active WordPress sites.

    Args:
        index: Index object for wp-rules
        is_updated: Whether files were actually updated
    """
    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        logger.info(
            "wordpress security plugin not enabled, skipping wp-rules"
            " deployment"
        )
        return

    if not is_updated:
        logger.info("wp-rules not updated, skipping deployment")
        return

    logger.info("Starting wp-rules deployment to WordPress sites")

    wp_rules_data = get_updated_wp_rules_data(index)
    if not wp_rules_data:
        logger.error("No valid wp-rules found, skipping deployment")
        return

    # Get version and create ruleset dict with version and rules
    wp_rules_version = get_wp_ruleset_version(index)
    ruleset_dict = {
        "version": wp_rules_version,
        "rules": wp_rules_data,
    }
    wp_rules_php = format_php_with_embedded_json(ruleset_dict)

    await _deploy_wp_rules_php(wp_rules_php)


_redeploy_rules_php_lock = asyncio.Lock()
# Separate flag is needed because lock.locked() is always True inside the
# holder's context, so it cannot indicate whether another caller coalesced.
_redeploy_rules_php_pending = False


async def redeploy_rules_php() -> None:
    """
    Re-deploy rules.php to all WordPress sites.

    Used when globally disabled rules change, requiring rules.php
    to be regenerated with updated rule filtering.

    Uses a coalescing lock: if a redeployment is already running,
    the request is merged into the current run rather than starting
    a duplicate deployment.
    """
    global _redeploy_rules_php_pending

    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        logger.info(
            "wordpress security plugin not enabled, skipping wp-rules"
            " redeployment"
        )
        return

    if _redeploy_rules_php_lock.locked():
        _redeploy_rules_php_pending = True
        logger.info("wp-rules redeployment already in progress, coalescing")
        return

    async with _redeploy_rules_php_lock:
        while True:
            _redeploy_rules_php_pending = False

            logger.info(
                "Starting wp-rules redeployment (global disable change)"
            )

            wp_rules_php = await load_wp_rules_php()
            if not wp_rules_php:
                logger.warning("Could not load wp-rules for redeployment")
                return

            await _deploy_wp_rules_php(wp_rules_php)

            if not _redeploy_rules_php_pending:
                break
            logger.info("Re-running wp-rules redeployment (coalesced request)")


async def redeploy_waf_for_all_sites() -> None:
    """Global WAF turn-on: deploy rules.php and disabled-rules.php (stamping
    disabled_rules_sync_ts) to all sites, matching install-with-WAF-on.

    Wraps rather than extends redeploy_rules_php, which is also the
    global-rule-change path where restamping sync_ts would skip unconsumed
    changelog actions.
    """
    await redeploy_rules_php()
    await update_disabled_rules_on_sites()


def _data_dir_belongs_to_user(
    dir_fd: int, data_dir, docroot: str, uid: int
) -> bool:
    """Confirm the data dir is the site owner's, checked on the fd.

    get_data_dir honours the owner-controlled WP_CONTENT_DIR, so an
    unvalidated path can lead into another tenant's account.
    """
    dir_uid = os.fstat(dir_fd).st_uid
    if dir_uid == uid:
        return True
    logger.warning(
        "Skipping WAF file removal: data dir %s for %s is owned by"
        " uid %s, expected %s",
        data_dir,
        docroot,
        dir_uid,
        uid,
    )
    return False


def _waf_confirmed_disabled(username: str, data_dir) -> bool:
    """Confirm WAF is still off for *username* just before unlinking.

    Checked here rather than by the caller because the read and the unlink
    share this worker thread with no await between them, which narrows the
    turn-on race window — it does not close it, since the event loop keeps
    running and a turn-on can still land between the read and the unlink. An
    unreadable status keeps the files.
    """
    try:
        if _is_waf_enabled_for_user_sync(username):
            logger.info(
                "WAF enabled for user %s, keeping WAF files in %s",
                username,
                data_dir,
            )
            return False
    except Exception:
        logger.warning(
            "Could not confirm WAF status for user %s, keeping WAF files"
            " in %s",
            username,
            data_dir,
            exc_info=True,
        )
        return False
    return True


def _remove_waf_files_for_dir(
    data_dir, docroot: str, uid: int, username: str | None = None
) -> bool:
    """Remove WAF files from data_dir via a symlink-safe dir fd.

    Returns True when the site reached the WAF-off target — files removed or
    already absent — and False when a guard or error left it unconverged.
    """
    try:
        dir_fd = open_dir_no_symlinks(data_dir)
    except FileNotFoundError:
        return True
    except OSError as exc:
        if exc.errno in (errno.ELOOP, errno.ENOTDIR):
            logger.warning(
                "Skipping WAF file removal: data dir %s is a symlink",
                data_dir,
            )
            return False
        logger.error("Failed to open data dir for %s: %s", docroot, exc)
        return False
    try:
        if not _data_dir_belongs_to_user(dir_fd, data_dir, docroot, uid):
            return False
        if username is not None and not _waf_confirmed_disabled(
            username, data_dir
        ):
            return False

        converged = True
        for filename in ("rules.php", "disabled-rules.php"):
            try:
                os.remove(filename, dir_fd=dir_fd)
                logger.info(
                    "Removed %s from %s (WAF disabled)", filename, docroot
                )
            except FileNotFoundError:
                pass
            except OSError as e:
                converged = False
                logger.error(
                    "Failed to remove %s from %s: %s", filename, docroot, e
                )
        return converged
    finally:
        os.close(dir_fd)


async def _remove_waf_files_from_sites(
    sites: list[WPSite], username: str | None = None
) -> list[str]:
    """Remove WAF files (rules.php, disabled-rules.php) from the given sites.

    Deletion goes through open_dir_no_symlinks + dir_fd so a site owner
    cannot redirect the root agent's removal via a symlinked data dir, the
    same symlink-safe pattern as delete_plugin_files.

    *username*, when given, is confirmed again immediately before each unlink
    — resolving the data dir yields, and a WAF turn-on landing in that window
    would otherwise lose its freshly deployed ruleset.

    Returns the docroots left unconverged — a guard or error kept the files.
    """
    unconverged = []
    for site in sites:
        try:
            data_dir = await cli.get_data_dir(site)
        except Exception as e:
            logger.error(
                "Failed to resolve data dir for %s: %s", site.docroot, e
            )
            unconverged.append(site.docroot)
            continue
        if not await asyncio.to_thread(
            _remove_waf_files_for_dir,
            data_dir,
            site.docroot,
            site.uid,
            username,
        ):
            unconverged.append(site.docroot)
    return unconverged


async def redeploy_waf_for_user(username: str, sink=None) -> None:
    """Deploy rules.php and disabled-rules.php for a user's sites (WAF turn-on).

    Caller must have confirmed WAF is enabled for the user. disabled-rules.php
    is deployed even if the ruleset fails to load, because it stamps
    disabled_rules_sync_ts — matching install, which writes it whenever WAF is
    on regardless of rules.php content.
    """
    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        return

    loop = asyncio.get_running_loop()
    try:
        user_info = await loop.run_in_executor(None, pwd.getpwnam, username)
    except KeyError:
        logger.warning(
            "User %s not found, skipping WAF rules redeploy", username
        )
        return

    sites = await loop.run_in_executor(None, get_installed_sites)
    user_sites = [s for s in sites if s.uid == user_info.pw_uid]
    if not user_sites:
        return

    live_sites = []
    for site in user_sites:
        if await remove_site_if_missing(sink, site):
            continue
        live_sites.append(site)
    if not live_sites:
        return

    updated = set()
    failed = set()
    wp_rules_php = await load_wp_rules_php()
    if wp_rules_php:
        for site in live_sites:
            await update_wp_rules_for_site(
                site, user_info, wp_rules_php, updated, failed
            )
    else:
        logger.warning("Could not load wp-rules for user redeploy")

    disabled_rules_ts = time.time()
    dr_updated = set()
    dr_failed = set()
    for site in live_sites:
        await update_disabled_rules_for_site(
            site, user_info, disabled_rules_ts, dr_updated, dr_failed
        )

    logger.info(
        "Redeployed WAF artifacts for user %s: rules %d ok/%d failed,"
        " disabled-rules %d ok/%d failed",
        username,
        len(updated),
        len(failed),
        len(dr_updated),
        len(dr_failed),
    )


async def remove_waf_rules_for_user(username: str) -> None:
    """Remove WAF files from all sites belonging to a user."""
    loop = asyncio.get_running_loop()
    try:
        user_info = await loop.run_in_executor(None, pwd.getpwnam, username)
    except KeyError:
        logger.warning(
            "User %s not found, skipping WAF rules removal", username
        )
        return

    sites = await loop.run_in_executor(None, get_installed_sites)
    user_sites = [s for s in sites if s.uid == user_info.pw_uid]
    await _remove_waf_files_from_sites(user_sites, username)


async def remove_waf_rules_for_all_sites() -> None:
    """Remove WAF files from all installed sites (global WAF disable)."""
    loop = asyncio.get_running_loop()
    sites = await loop.run_in_executor(None, get_installed_sites)
    logger.info(
        "Global WAF disabled, removing WAF files from %d site(s)",
        len(sites),
    )
    await _remove_waf_files_from_sites(sites)


async def remove_waf_rules_for_waf_disabled_users() -> None:
    """Remove WAF files from every installed site whose owner is WAF-off."""
    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        # remove_all_installed owns cleanup when the plugin is off.
        return

    loop = asyncio.get_running_loop()
    sites = await loop.run_in_executor(None, get_installed_sites)
    if not sites:
        return

    clear_caches()

    sites_by_user = defaultdict(list)
    for site in sites:
        sites_by_user[site.uid].append(site)

    with inactivity.track.task("wp-waf-off-removal"):
        waf_off_users = 0
        users_converged = 0
        sites_unconverged = []
        for uid, user_sites in sites_by_user.items():
            try:
                user_info = await loop.run_in_executor(None, pwd.getpwuid, uid)
                username = user_info.pw_name
            except Exception as error:
                logger.warning(
                    "Could not resolve uid %s, leaving WAF files on %d"
                    " site(s). Reason: %s",
                    uid,
                    len(user_sites),
                    error,
                )
                sites_unconverged.extend(site.docroot for site in user_sites)
                continue

            try:
                waf_enabled = await is_waf_enabled_for_user(username)
            except Exception:
                logger.warning(
                    "Could not check WAF status for user %s, leaving WAF"
                    " files on %d site(s)",
                    username,
                    len(user_sites),
                    exc_info=True,
                )
                sites_unconverged.extend(site.docroot for site in user_sites)
                continue

            if waf_enabled:
                continue

            waf_off_users += 1
            unconverged = await _remove_waf_files_from_sites(
                user_sites, username
            )
            if unconverged:
                sites_unconverged.extend(unconverged)
            else:
                users_converged += 1

        if sites_unconverged:
            logger.warning(
                "WAF-off removal pass left %d site(s) unconverged: %s",
                len(sites_unconverged),
                ", ".join(sites_unconverged),
            )
        logger.info(
            "WAF-off removal pass complete: %d of %d WAF-off user(s)"
            " converged, %d site(s) checked, %d site(s) unconverged",
            users_converged,
            waf_off_users,
            len(sites),
            len(sites_unconverged),
        )


async def apply_waf_default_change(sink=None) -> None:
    """Redeploy/remove rules.php for users without an explicit waf_enabled override."""
    global _apply_waf_default_pending

    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        return

    if _apply_waf_default_lock.locked():
        _apply_waf_default_pending = True
        logger.info("waf_default change already in progress, coalescing")
        return

    async with _apply_waf_default_lock:
        while True:
            _apply_waf_default_pending = False

            if not _get_global_waf_enabled():
                return

            new_default = _get_waf_default()
            usernames = await hosting_panel.HostingPanel().get_users()
            loop = asyncio.get_running_loop()

            for username in usernames:
                try:
                    if await loop.run_in_executor(
                        None,
                        _user_has_explicit_waf_override_sync,
                        username,
                    ):
                        continue
                    if new_default:
                        await redeploy_waf_for_user(username, sink=sink)
                    else:
                        await remove_waf_rules_for_user(username)
                except Exception as e:
                    logger.warning(
                        "Failed to apply waf_default change for user %s: %s",
                        username,
                        e,
                    )

            if not _apply_waf_default_pending:
                break
            logger.info("Re-running waf_default change (coalesced request)")


async def redeploy_ai_bot_for_user(username: str, sink=None) -> None:
    """Rewrite plugin_config.php for a user's sites (per-account ai-bot change).

    Unlike WAF there is no rules file to add or remove — plugin_config.php
    always exists and carries ai_bot_protection + preset as fields, so a
    per-account change is a rewrite in either direction.
    """
    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        return
    loop = asyncio.get_running_loop()
    try:
        user_info = await loop.run_in_executor(None, pwd.getpwnam, username)
    except KeyError:
        logger.warning(
            "User %s not found, skipping ai-bot config redeploy", username
        )
        return
    sites = await loop.run_in_executor(None, get_installed_sites)
    user_sites = [s for s in sites if s.uid == user_info.pw_uid]
    if not user_sites:
        return
    await update_plugin_config_on_sites(user_sites, sink=sink)


async def apply_ai_bot_default_change(sink=None) -> None:
    """Rewrite plugin_config.php for users without an explicit ai_bot_protection override."""
    global _apply_ai_bot_default_pending

    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        return

    if _apply_ai_bot_default_lock.locked():
        _apply_ai_bot_default_pending = True
        logger.info(
            "ai_bot_protection_default change already in progress, coalescing"
        )
        return

    async with _apply_ai_bot_default_lock:
        while True:
            _apply_ai_bot_default_pending = False

            if not _get_global_ai_bot_protection():
                return

            usernames = await hosting_panel.HostingPanel().get_users()
            loop = asyncio.get_running_loop()

            # One site fetch + one bulk write: per-user redeploys would
            # refetch the full site list for every user, blocking the
            # message sink for minutes on large servers. The bulk writer
            # resolves each owner's effective values at write time and
            # isolates per-user failures. Catching here keeps a transient
            # failure from dropping a coalesced re-run.
            try:
                uids = await loop.run_in_executor(
                    None, _no_override_uids_sync, usernames
                )
                sites = await loop.run_in_executor(None, get_installed_sites)
                await update_plugin_config_on_sites(
                    [site for site in sites if site.uid in uids], sink=sink
                )
            except Exception as e:
                logger.warning(
                    "Failed to apply ai_bot_protection_default change: %s", e
                )

            if not _apply_ai_bot_default_pending:
                break
            logger.info(
                "Re-running ai_bot_protection_default change"
                " (coalesced request)"
            )


def generate_disabled_rules_php(domain: str, timestamp: float) -> str:
    """
    Generate the disabled-rules.php content for a specific domain.

    Only includes domain-specific disabled rules. Globally disabled rules
    are handled separately by filtering them out of rules.php.

    Args:
        domain: The domain to generate disabled rules for
        timestamp: Unix timestamp to embed in the file

    Returns:
        PHP file content string
    """
    disabled_rule_ids = WPDisabledRule.get_domain_disabled(
        domain, include_global=False
    )
    data = {
        "ts": timestamp,
        "rules": sorted(disabled_rule_ids),
    }
    return format_php_with_embedded_json(data)


async def update_disabled_rules_for_site(
    site: WPSite,
    user_info: pwd.struct_passwd,
    timestamp: float,
    updated: set,
    failed: set,
) -> None:
    """
    Deploy disabled-rules.php to a single WordPress site and track the result.

    Args:
        site: WordPress site to deploy to
        user_info: User information from pwd
        timestamp: Unix timestamp for both file content and DB record
        updated: Set to add site to if successful
        failed: Set to add site to if failed
    """
    gid = user_info.pw_gid

    try:
        data_dir = await ensure_site_data_directory(site, user_info)
        if data_dir is None:
            return
        disabled_rules_path = data_dir / "disabled-rules.php"
        php_content = generate_disabled_rules_php(site.domain, timestamp)
        write_plugin_data_file_atomically(
            disabled_rules_path, php_content, uid=site.uid, gid=gid
        )
        WordpressSite.update(disabled_rules_sync_ts=timestamp).where(
            WordpressSite.docroot == site.docroot
        ).execute()
        updated.add(site)
        logger.info("Updated disabled-rules for site %s", site.docroot)
    except Exception as error:
        failed.add(site)
        logger.error(
            "Failed to update disabled-rules for site %s: %s",
            site.docroot,
            error,
        )


async def update_disabled_rules_on_sites(
    domains: list[str] | None = None,
    sink=None,
) -> None:
    """
    Deploy disabled-rules.php to WordPress sites.

    If domains are specified, only updates sites for those domains.
    If domains is None, updates all installed sites (e.g., after a global
    disable/enable).

    Args:
        domains: List of domains to update, or None for all sites
        sink: Optional telemetry sink for remove_site_if_missing
    """
    if not Wordpress.SECURITY_PLUGIN_ENABLED:
        logger.info(
            "wordpress security plugin not enabled, skipping disabled-rules"
            " deployment"
        )
        return

    logger.info("Starting disabled-rules deployment to WordPress sites")

    clear_caches()

    if domains:
        sites = get_installed_sites_by_domains(domains)
    else:
        sites = get_installed_sites()

    if not sites:
        logger.info("No WordPress sites found for disabled-rules deployment")
        return

    def make_task(site, user_info, updated, failed):
        return update_disabled_rules_for_site(
            site, user_info, time.time(), updated, failed
        )

    await _deploy_to_sites(
        sites,
        make_task,
        task_name="disabled-rules",
        fingerprint="disabled-rules-update-skip-user",
        sink=sink,
    )


async def update_auth_everywhere(sink=None):
    """Update auth.php files for all existing WordPress sites."""
    logger.info("Updating auth.php files for existing WordPress sites")

    updated = set()
    failed = set()

    with inactivity.track.task("wp-auth-update"):
        try:
            clear_caches()

            # Get all installed sites from db
            installed_sites = get_installed_sites()

            if not installed_sites:
                logger.info("No installed WordPress sites found")
                return

            sites_by_user = defaultdict(list)
            for site in installed_sites:
                sites_by_user[site.uid].append(site)

            # Process users concurrently
            tasks = []
            for uid, sites in sites_by_user.items():
                try:
                    user_info = pwd.getpwuid(uid)
                except Exception as error:
                    log_message(
                        "Skipping auth update for WordPress sites on"
                        " {count} site(s) because they belong to user"
                        " {user} and it is not possible to retrieve"
                        " username for this user. Reason: {reason}",
                        format_args={
                            "count": len(sites),
                            "user": uid,
                            "reason": error,
                        },
                        level="warning",
                        component="wordpress",
                        fingerprint="wp-plugin-auth-update-skip-user",
                    )
                    continue

                for site in sites:
                    if await remove_site_if_missing(sink, site):
                        continue
                    task = update_site_auth(site, user_info, updated, failed)
                    tasks.append(task)

            # Run all site updates concurrently with a reasonable limit
            # Adjust max_concurrent based on your system's I/O capacity
            max_concurrent = 10
            for i in range(0, len(tasks), max_concurrent):
                batch = tasks[i : i + max_concurrent]
                await asyncio.gather(*batch, return_exceptions=True)

            logger.info(
                "Updated auth.php files for %d WordPress sites, %d failed",
                len(updated),
                len(failed),
            )

        except asyncio.CancelledError:
            logger.info(
                "Auth update for WordPress sites was cancelled. Auth was"
                " updated for %d sites",
                len(updated),
            )
        except Exception as error:
            logger.error("Error occurred during auth update. error=%s", error)
            raise


async def update_site_auth(site, user_info, updated, failed):
    """Process authentication setup for a single site."""
    try:
        await setup_site_authentication(site, user_info)
        updated.add(site)
    except Exception as error:
        failed.add(site)
        logger.error(
            "Failed to update auth for site=%s error=%s",
            site,
            error,
        )


def _is_docroot_present(site: WPSite) -> bool:
    # os.path.isdir() is not used here: it reports every stat failure as
    # "not a directory", so an unreadable docroot (EACCES/ESTALE/EIO) would
    # read as a deleted site and drop its DB row. Only ENOENT/ENOTDIR is gone.
    try:
        docroot_stat = os.stat(site.docroot)
    except (FileNotFoundError, NotADirectoryError):
        return False
    except (OSError, ValueError) as error:
        logger.warning(
            "Failed to stat %s, assuming the docroot is still present: %s",
            site.docroot,
            error,
        )
        return True
    return stat.S_ISDIR(docroot_stat.st_mode)


async def remove_site_if_missing(sink, site: WPSite) -> bool:
    """
    Remove the site from the local DB (and emit 'site_removed' on success) only
    when its docroot is confirmed gone (ENOENT/ENOTDIR); any other stat failure
    keeps the site. Returns True if the site was removed, False otherwise.
    """
    if _is_docroot_present(site):
        return False

    # Attempt to delete the site from the database first
    rows_deleted = delete_site(site)

    # Only send telemetry if the deletion was successful (at least one row was deleted)
    if rows_deleted > 0:
        if sink is not None:
            await telemetry.send_event(
                sink=sink,
                event="site_removed",
                site=site,
                version=site.version,
            )
    else:
        logger.warning(
            "Failed to delete missing site %s from database, no rows affected",
            site,
        )

        log_message(
            "Failed to delete missing site {site} from database",
            format_args={"site": site},
            level="warning",
            component="wordpress",
            fingerprint="wp-plugin-site-delete-failed",
        )

    return True


async def fix_site_data_file_permissions(
    site: WPSite, file_permissions: int
) -> bool:
    """Fix data file permissions for a single WordPress site."""
    try:
        data_dir = await cli.get_data_dir(site)

        try:
            dir_fd = open_dir_no_symlinks(data_dir)
        except OSError as exc:
            if exc.errno in (errno.ENOENT, errno.ELOOP, errno.ENOTDIR):
                return False
            raise

        try:
            # Fix directory permissions via fd.
            current_dir_mode = os.stat(dir_fd).st_mode & 0o777
            if current_dir_mode != 0o750:
                os.chmod(dir_fd, 0o750)

            for file_name in [
                "scan_data.php",
                "plugin_config.php",
                "auth.php",
                "rules.php",
                "disabled-rules.php",
                IP_WHITELIST_FILENAME,
            ]:
                try:
                    with open_nofollow(file_name, dir_fd=dir_fd) as file_fd:
                        st = os.fstat(file_fd)
                        if st.st_mode & 0o777 != file_permissions:
                            os.chmod(file_fd, file_permissions)
                except FileNotFoundError:
                    continue
                except OSError as exc:
                    if exc.errno == errno.ELOOP:
                        logger.warning(
                            "Skipping chmod: %s/%s is a symlink",
                            data_dir,
                            file_name,
                        )
                        continue
                    raise
        finally:
            os.close(dir_fd)

        return True
    except Exception as error:
        logger.error(
            "Failed to fix permissions for site=%s error=%s",
            site,
            error,
        )
        return False


async def fix_data_file_permissions_everywhere(sink):
    """
    Fix data file permissions for all WordPress sites with imunify-security plugin installed.

    Args:
        sink: The telemetry/event sink
    """
    fixed = set()
    failed = set()

    with inactivity.track.task("wp-plugin-fix-permissions"):
        try:
            clear_caches()

            # Get all installed sites
            installed_sites = get_installed_sites()
            if not installed_sites:
                return

            # Determine file permissions based on hosting panel
            from defence360agent.subsys.panels.hosting_panel import (
                HostingPanel,
            )
            from defence360agent.subsys.panels.plesk import Plesk

            file_permissions = (
                0o440 if HostingPanel().NAME == Plesk.NAME else 0o400
            )

            # Process sites
            for site in installed_sites:
                if await remove_site_if_missing(sink, site):
                    continue

                success = await fix_site_data_file_permissions(
                    site, file_permissions
                )
                if success:
                    fixed.add(site)
                else:
                    failed.add(site)

            logger.info(
                "Fixed data file permissions for %d WordPress sites, %d"
                " failed",
                len(fixed),
                len(failed),
            )

        except asyncio.CancelledError:
            logger.info(
                "Fixing data file permissions was cancelled. Permissions were"
                " fixed for %d sites",
                len(fixed),
            )
        except Exception as error:
            logger.error(
                "Error occurred during permission fixing. error=%s", error
            )


class WordPressSiteInstaller:
    """
    Handles installation of imunify-security plugin on WordPress sites.

    This class processes WordPress sites and installs the imunify-security
    plugin, including setting up authentication, scan data files, and rules.
    """

    install_plugin = True
    telemetry_event = "installed_by_imunify"
    task_name = "wp-plugin-installation"
    log_fingerprint_skip_user = "wp-plugin-install-skip-user"
    messages = {
        "start": "Installing imunify-security wp plugin",
        "complete": "Installed imunify-security wp plugin on {count} sites",
        "found": "Found {count} site(s) for installation",
        "error": "Failed to install plugin to site={site} error={error}",
        "cancelled": (
            "Installation of imunify-security wp plugin was cancelled. "
            "Plugin was installed for {count} sites"
        ),
        "exception": (
            "Error occurred during plugin installation. error={error}"
        ),
        "skip_user": (
            "Skipping installation of WordPress plugin on "
            "{count} site(s) because they belong to user "
            "{user} and it is not possible to retrieve "
            "username for this user. Reason: {reason}"
        ),
    }

    def __init__(self, sink, sites):
        self.sink = sink
        self.sites = sites
        self.processed = set()
        self.authenticated = set()
        self.rules_installed = set()
        self.failed_rules_updates = set()
        self.disabled_rules_installed = set()
        self.failed_disabled_rules_updates = set()
        self.disabled_rules_ts: float | None = None
        self.failed_auth = set()
        self._current_site: WPSite | None = None

    async def is_site_ready(self, site):
        """
        Check if site is ready for processing.

        Override in subclasses to implement different readiness checks.

        Args:
            site: The WordPress site to check.

        Returns:
            bool: True if the site is ready for processing, False otherwise.
        """
        is_wordpress_installed = await cli.is_wordpress_installed(site)
        if not is_wordpress_installed:
            logger.warning(
                "WordPress site is not accessible using WP CLI. site=%s",
                site,
            )
            log_message(
                "WordPress site is not accessible using WP CLI. site={site}",
                format_args={"site": site},
                level="warning",
                component="wordpress",
                fingerprint="wp-plugin-cli-not-accessible",
            )
            return False
        return True

    def _record_processed_site(self, site, version):
        """
        Record a successfully processed site and persist it to the database.

        Each site is inserted immediately so that it is tracked in the DB
        at all times — even if the overall installation loop is cancelled
        mid-run.

        Override in subclasses to implement different recording logic.

        Args:
            site: The WordPress site that was processed.
            version: The plugin version installed on the site.
        """
        self.processed.add(site)
        insert_installed_sites({site})
        self._stamp_disabled_rules_sync_ts(site)

    async def _revert_in_flight_site(self):
        """Revert the site that was mid-processing when cancellation occurred.

        Deletes data files and, if this processor installs plugins,
        attempts to uninstall the partially-installed plugin.
        Each step runs independently so one failure doesn't skip the other.
        """
        site = self._current_site
        self._current_site = None
        try:
            await delete_plugin_files(site)
        except Exception as error:
            logger.warning(
                "Failed to delete data files for in-flight site %s: %s",
                site,
                error,
            )
        if self.install_plugin:
            await cli.try_plugin_uninstall(site)

    def _stamp_disabled_rules_sync_ts(self, site: WPSite) -> None:
        """
        Stamp disabled_rules_sync_ts for a single site after it has been
        inserted into the DB.

        Called from ``_record_processed_site`` so the DB row already exists.
        """
        if site not in self.disabled_rules_installed:
            return
        if self.disabled_rules_ts is None:
            return
        WordpressSite.update(
            disabled_rules_sync_ts=self.disabled_rules_ts
        ).where(WordpressSite.docroot == site.docroot).execute()

    async def run(self):
        """
        Process WordPress sites for imunify-security plugin operations.

        Returns:
            set: The set of successfully processed sites.
        """
        logger.info(self.messages["start"])
        telemetry_tasks = []

        with inactivity.track.task(self.task_name):
            try:
                clear_caches()

                if not self.sites:
                    logger.info("No WordPress sites found, nothing to do")
                    return self.processed

                logger.info(
                    self.messages["found"].format(count=len(self.sites))
                )

                # Create SystemConfig once for all users
                admin_config = SystemConfig()

                # Create wp rules once for all users
                wp_rules_php = await load_wp_rules_php()
                # Always set the timestamp, even when no disabled rules
                # currently exist: previously disabled-then-enabled rules
                # require deploying an empty disabled-rules.php.
                self.disabled_rules_ts = time.time()

                versions = await get_imunify_package_versions()

                # Group sites by user id
                sites_by_user = defaultdict(list)
                for site in self.sites:
                    sites_by_user[site.uid].append(site)

                # Now iterate over the grouped sites
                for uid, sites in sites_by_user.items():
                    try:
                        user_info = pwd.getpwuid(uid)
                        username = user_info.pw_name
                    except Exception as error:
                        log_message(
                            self.messages["skip_user"],
                            format_args={
                                "count": len(sites),
                                "user": uid,
                                "reason": error,
                            },
                            level="warning",
                            component="wordpress",
                            fingerprint=self.log_fingerprint_skip_user,
                        )
                        continue

                    (
                        last_scan_time,
                        next_scan_time,
                        malware_by_site,
                    ) = await _get_scan_data_for_user(
                        self.sink, user_info, admin_config
                    )
                    plugin_config = prepare_plugin_config(username)

                    for site in sites:
                        if await remove_site_if_missing(self.sink, site):
                            continue

                        try:
                            # Check if site is ready for processing (WP CLI accessible + other checks)
                            if not await self.is_site_ready(site):
                                continue

                            self._current_site = site

                            # Prepare scan data
                            scan_data = prepare_scan_data(
                                last_scan_time,
                                next_scan_time,
                                username,
                                site,
                                malware_by_site,
                                versions=versions,
                            )

                            # Resolve the data directory once; the scan-data
                            # and plugin-config writes below reuse it.
                            data_dir = await ensure_site_data_directory(
                                site, user_info
                            )
                            if data_dir is None:
                                continue

                            # Create data files (scan data, plugin config, auth token)
                            await update_scan_data_file(
                                site,
                                scan_data,
                                user_info=user_info,
                                data_dir=data_dir,
                            )
                            await update_plugin_config_file(
                                site,
                                plugin_config,
                                user_info=user_info,
                                data_dir=data_dir,
                            )
                            ip_whitelist = await _load_ip_whitelist_or_skip()
                            if ip_whitelist is not None:
                                text, _ = ip_whitelist
                                await _sync_ip_whitelist_or_skip(
                                    site,
                                    text,
                                    user_info=user_info,
                                    data_dir=data_dir,
                                )
                            await update_site_auth(
                                site,
                                user_info,
                                self.authenticated,
                                self.failed_auth,
                            )

                            if await _waf_enabled_or_default(username):
                                if wp_rules_php:
                                    await update_wp_rules_for_site(
                                        site,
                                        user_info,
                                        wp_rules_php,
                                        self.rules_installed,
                                        self.failed_rules_updates,
                                    )
                                await update_disabled_rules_for_site(
                                    site,
                                    user_info,
                                    self.disabled_rules_ts,
                                    self.disabled_rules_installed,
                                    self.failed_disabled_rules_updates,
                                )
                            else:
                                await asyncio.to_thread(
                                    _remove_waf_files_for_dir,
                                    data_dir,
                                    site.docroot,
                                    site.uid,
                                    username,
                                )

                            # Install the plugin
                            if self.install_plugin:
                                await cli.plugin_install(site)

                            # Get the version of the plugin
                            version = await cli.get_plugin_version(site)
                            if version:
                                site = WPSite.build_with_version(site, version)

                            # Record the processed site
                            self._record_processed_site(site, version)

                            self._current_site = None

                            telemetry_tasks.append(
                                asyncio.create_task(
                                    telemetry.send_event(
                                        sink=self.sink,
                                        event=self.telemetry_event,
                                        site=site,
                                        version=version,
                                    )
                                )
                            )
                        except Exception as error:
                            self._current_site = None
                            logger.error(
                                self.messages["error"].format(
                                    site=site, error=repr(error)
                                )
                            )
                logger.info(
                    self.messages["complete"].format(count=len(self.processed))
                )
                if self.failed_auth:
                    logger.warning(
                        "Failed to authenticate %d sites",
                        len(self.failed_auth),
                    )
                if self.failed_rules_updates:
                    logger.warning(
                        "Failed to install wp-rules on %d sites",
                        len(self.failed_rules_updates),
                    )
                if self.failed_disabled_rules_updates:
                    logger.warning(
                        "Failed to install disabled-rules on %d sites",
                        len(self.failed_disabled_rules_updates),
                    )

            except asyncio.CancelledError:
                if self._current_site:
                    await self._revert_in_flight_site()
                logger.info(
                    self.messages["cancelled"].format(
                        count=len(self.processed)
                    )
                )
            except Exception as error:
                logger.error(
                    self.messages["exception"].format(error=repr(error))
                )
                raise
            finally:
                if telemetry_tasks:
                    results = await asyncio.gather(
                        *telemetry_tasks, return_exceptions=True
                    )
                    for result in results:
                        if isinstance(result, Exception):
                            logger.warning(
                                "Failed to send telemetry: %s", result
                            )

        return self.processed


class WordPressSiteAdopter(WordPressSiteInstaller):
    """
    Handles adoption of existing WordPress sites with imunify-security plugin.

    Adoption is a special case of installation where the site already has
    the plugin installed but is not tracked in our database.
    """

    install_plugin = False
    telemetry_event = "site_found"
    task_name = "wp-plugin-adoption"
    log_fingerprint_skip_user = "wp-plugin-adopt-skip-user"
    messages = {
        "start": "Adopting imunify-security wp plugin",
        "complete": "Adopted imunify-security wp plugin on {count} sites",
        "found": "Found {count} site(s) for adoption",
        "error": "Failed to adopt plugin to site={site} error={error}",
        "cancelled": (
            "Adoption of imunify-security wp plugin was cancelled. "
            "Plugin was adopted for {count} sites"
        ),
        "exception": "Error occurred during plugin adoption. error={error}",
        "skip_user": (
            "Skipping adoption of WordPress plugin on "
            "{count} site(s) because they belong to user "
            "{user} and it is not possible to retrieve "
            "username for this user. Reason: {reason}"
        ),
    }

    def __init__(self, sink, sites):
        super().__init__(sink, sites)
        # Load existing docroots from database for adoption logic
        self.existing_docroots = {
            r.docroot for r in WordpressSite.select(WordpressSite.docroot)
        }

    def _record_processed_site(self, site, version):
        """
        Record a successfully adopted site and persist it immediately.

        For adoption, sites that already exist in the database (flagged as
        manually deleted) have their flag cleared. New sites are inserted
        into the database right away.

        Args:
            site: The WordPress site that was processed.
            version: The plugin version installed on the site.
        """
        self.processed.add(site)
        if site.docroot in self.existing_docroots:
            # Site exists in DB but is flagged - clear flag
            clear_manually_deleted_flag(site)
            update_site_identity(site)
            if version:
                update_site_version(site, version)
        else:
            insert_installed_sites({site})
        self._stamp_disabled_rules_sync_ts(site)

    async def is_site_ready(self, site):
        """
        Check if site is ready for adoption.

        Args:
            site: The WordPress site to check.

        Returns:
            bool: True if the site is ready for adoption, False otherwise.
        """
        if not await super().is_site_ready(site):
            return False

        # Verify plugin is actually installed
        is_installed = await cli.is_plugin_installed(site)
        if not is_installed:
            logger.warning(
                "Plugin not installed on site %s, skipping adoption",
                site,
            )
            return False

        return True
