package ru.kirillius.pf.sdn.service;

import lombok.Getter;
import org.json.JSONObject;
import org.json.JSONTokener;
import org.springframework.stereotype.Service;
import ru.kirillius.java.utils.events.EventListener;
import ru.kirillius.json.JSONUtility;
import ru.kirillius.pf.sdn.api.Networking.IPv4Subnet;
import ru.kirillius.pf.sdn.api.Networking.NetworkScope;
import ru.kirillius.pf.sdn.dto.ResolverCacheEntry;
import ru.kirillius.pf.sdn.core.AppService;
import ru.kirillius.pf.sdn.core.Context;
import ru.kirillius.pf.sdn.core.ContextEventsHandler;
import ru.kirillius.pf.sdn.core.Subscription.SubscriptionService;
import ru.kirillius.pf.sdn.core.Util.IPv4Util;
import ru.kirillius.utils.logging.SystemLogger;

import java.io.File;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
import java.time.Instant;
import java.time.temporal.ChronoUnit;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;

/**
 * Builds the effective set of network resources by combining subscriptions, caches, and filters.
 */
@Service

public class NetworkingService extends AppService {
    private final static String CTX = NetworkingService.class.getSimpleName();
    private final ExecutorService executor = Executors.newSingleThreadExecutor();
    private final File domainCacheFile;
    private final File asCacheFile;
    private final EventListener<NetworkResourceConfig> resourceUpdateSubscription;
    private final EventListener<ContextEventsHandler.ConfigChangeContext> configChangeSubscription;
    private final AtomicReference<Future<?>> updateProcess = new AtomicReference<>();
    @Getter
    private final NetworkResourceConfig inputResources = new NetworkResourceConfig();
    @Getter
    private final NetworkResourceConfig outputResources = new NetworkResourceConfig();

    private final Map<Integer, List<IPv4Subnet>> prefixCache = new ConcurrentHashMap<>();
    private final Map<String, ResolverCacheEntry> domainCache = new ConcurrentHashMap<>();

    /**
     * Creates the networking service, wiring subscriptions and restoring cached state.
     */
    public NetworkingService(Context context) {
        super(context);
        inputResources.clear();
        inputResources.add(context.getConfig().getCustomResources());
        resourceUpdateSubscription = context.getEventsHandler().getSubscriptionsUpdateEvent().add(bundle -> rebuildInputs());
        configChangeSubscription = context.getEventsHandler().getConfigChangeEvent().add(changeContext -> {
            var filtersChanges = !changeContext.getCurrent().getFilteredResources().equals(changeContext.getInitial().getFilteredResources());
            var resChanges = !changeContext.getCurrent().getCustomResources().equals(changeContext.getInitial().getCustomResources());
            if (resChanges || filtersChanges) {
                NetworkingService.this.rebuildInputs();
            }
        });
        domainCacheFile = new File(context.getConfig().getCacheDirectory(), "domain-cache.json");
        asCacheFile = new File(context.getConfig().getCacheDirectory(), "as-cache.json");
        if (asCacheFile.exists() && context.getConfig().isCachingAS()) {
            SystemLogger.message("Loading as cache file", CTX);
            try (var is = new FileInputStream(asCacheFile)) {
                var json = new JSONObject(new JSONTokener(is));
                json.keySet().forEach(key -> {
                    var as = Integer.parseInt(key);
                    prefixCache.put(as, JSONUtility.deserializeCollection(json.getJSONArray(key), IPv4Subnet.class, null).toList());
                });
            } catch (Exception e) {
                SystemLogger.error("Failed to load as cache file " + asCacheFile.getPath(), CTX, e);
            }
        }

        if (domainCacheFile.exists() && context.getConfig().isCachingAS()) {
            SystemLogger.message("Loading domain cache file", CTX);
            try (var is = new FileInputStream(domainCacheFile)) {
                var json = new JSONObject(new JSONTokener(is));
                json.keySet().forEach(host -> {
                    domainCache.put(host, JSONUtility.deserializeStructure(json.getJSONObject(host), ResolverCacheEntry.class));
                });
            } catch (Exception e) {
                SystemLogger.error("Failed to load domain cache file " + asCacheFile.getPath(), CTX, e);
            }
        }


    }

    public void performAutoresolve() {
        var current = new HashSet<IPv4Subnet>();
        domainCache.forEach((host, entry) -> current.addAll(entry.getAddresses().keySet()));

        resolveDomains(List.copyOf(context.getServiceManager().getService(SubscriptionService.class).getAutoResolvingDomains()));

        var resolved = new HashSet<IPv4Subnet>();
        domainCache.forEach((host, entry) -> resolved.addAll(entry.getAddresses().keySet()));

        if (resolved.size() != current.size()) {
            rebuildInputs();
        } else {
            var updated = false;
            for (var subnet : resolved) {
                if (!current.contains(subnet)) {
                    updated = true;
                    break;
                }
            }
            if (updated) {
                rebuildInputs();
            }
        }

        if(!context.getConfig().isCachingDomains()){
            domainCache.clear();
        }
    }

    private void rebuildInputs() {
        inputResources.clear();
        inputResources.add(context.getConfig().getCustomResources());
        inputResources.add(context.getServiceManager().getService(SubscriptionService.class).getOutputResources());
        triggerUpdate(false);
    }

    /**
     * Indicates whether an update job is currently executing.
     */
    public boolean isUpdatingNow() {
        var future = updateProcess.get();
        return future != null && !future.isDone() && !future.isCancelled();
    }

    /**
     * Schedules an update of network resources, optionally ignoring cached prefixes.
     */
    public void triggerUpdate(boolean ignoreCache) {
        if (isUpdatingNow()) {
            return;
        }
        SystemLogger.message("Updating network manager", CTX);

        updateProcess.set(executor.submit(() -> {
            try {
                SystemLogger.message("Update is started", CTX);
                var config = context.getConfig();
                var filteredResources = config.getFilteredResources();

                var domains = new HashSet<>(inputResources.getDomains());
                filteredResources.getDomains().forEach(domains::remove);
                //check domain overlaps

                var domainsToRemove = new HashSet<String>();
                for (var domainToMatch : domains) {
                    var pattern = "." + domainToMatch;
                    for (var domain : domains) {
                        if (domain.endsWith(pattern)) {
                            domainsToRemove.add(domain);
                        }
                    }
                }

                domains.removeAll(domainsToRemove);

                var asn = new ArrayList<>(inputResources.getASN());
                asn.removeAll(filteredResources.getASN());

                var asnToFetch = new ArrayList<>(asn);
                if (!ignoreCache) {
                    asnToFetch.removeAll(prefixCache.keySet());
                }

                fetchPrefixes(asnToFetch);

                if (config.isCachingAS()) {
                    try (var os = new FileOutputStream(asCacheFile)) {
                        var json = new JSONObject();
                        prefixCache.forEach((key, asnList) -> {
                            json.put(String.valueOf(key), JSONUtility.serializeCollection(asnList, IPv4Subnet.class, null));
                        });
                        os.write(json.toString().getBytes());
                    } catch (IOException e) {
                        SystemLogger.error("Unable to write file " + asCacheFile.getPath(), CTX, e);
                    }
                }

                resolveDomains(List.copyOf(context.getServiceManager().getService(SubscriptionService.class).getAutoResolvingDomains()));

                if (config.isCachingDomains()) {
                    try (var os = new FileOutputStream(domainCacheFile)) {
                        var json = new JSONObject();
                        domainCache.forEach((key, entry) -> {
                            var serialized = JSONUtility.serializeStructure(entry);
                            json.put(String.valueOf(key), serialized);
                        });
                        os.write(json.toString().getBytes());
                    } catch (IOException e) {
                        SystemLogger.error("Unable to write file " + domainCacheFile.getPath(), CTX, e);
                    }
                }

                var subnets = new HashSet<>(inputResources.getSubnets());
                asn.forEach(n -> {
                    var cached = prefixCache.get(n);
                    if (cached == null) {
                        return;
                    }
                    subnets.addAll(cached);
                    SystemLogger.message("Using " + cached.size() + " subnets from AS" + n, CTX);
                });

                //добавляем отрезолвенные домены
                domains.forEach(domain -> {
                    var entry = domainCache.get(domain);
                    if (entry != null) {
                        subnets.addAll(entry.getAddresses().keySet());
                    }
                });

                filteredResources.getSubnets().forEach(subnets::remove);

                SystemLogger.message("Trying to summary " + subnets.size() + " subnets...", CTX);

                var merged = IPv4Util.summarySubnets(subnets, config.getMergeSubnetsWithUsage());
                var unmerged = new AtomicInteger();
                subnets.forEach(subnet -> {
                    if (!merged.getMergedSubnets().contains(subnet)) {
                        unmerged.getAndIncrement();
                    }
                });

                SystemLogger.message(subnets.size() + " subnets has been summarized and merged to " + merged.getResult().size() + " new subnets. Unmerged: " + unmerged.get(), CTX);

                outputResources.setASN(Collections.unmodifiableList(asn));
                outputResources.setSubnets(merged.getResult());
                outputResources.setDomains(domains.stream().toList());

                SystemLogger.message("Update is complete", CTX);

                try {
                    context.getEventsHandler().getNetworkManagerUpdateEvent().invoke(outputResources);
                } catch (Exception e) {
                    SystemLogger.error("Unable to invoke update event", CTX, e);
                }
            } catch (Exception e) {
                SystemLogger.error("Something went wrong on update", CTX, e);
            }
        }));
    }

    private void resolveDomains(List<String> domains) {
        var resolvedSubnets = new ArrayList<IPv4Subnet>();
        var resolver = context.getServiceManager().getService(DomainResolverService.class);
        for (var domain : domains) {
            var task = resolver.resolve(domain);
            while (!task.isDone() && !task.isCancelled()) {
                Thread.yield();
            }
            try {
                var subnets = task.get();
                var entry = domainCache.get(domain);

                if(entry == null) {
                    entry = new ResolverCacheEntry();
                    domainCache.put(domain, entry);
                }

                var addresses = entry.getAddresses();
                entry.setLastUpdate(Instant.now());
                subnets.forEach(subnet -> addresses.put(subnet, Instant.now()));
                resolvedSubnets.addAll(domainCache.get(domain).getAddresses().keySet());
            } catch (InterruptedException | ExecutionException e) {
                SystemLogger.error("Error happened while resolving domain " + domain, CTX, e);
            }
        }

        //remove old entries
        for (var domain : domainCache.keySet()) {
            var entry = domainCache.get(domain);
            var addresses = entry.getAddresses();
            for (var subnet : addresses.keySet()) {
                var time = addresses.get(subnet);
                if (time.isBefore(Instant.now().minus(context.getConfig().getDomainsTimeToLive(), ChronoUnit.HOURS))) {
                    addresses.remove(subnet);
                }
            }
            if (addresses.isEmpty()) {
                domainCache.remove(domain);
            }
        }
    }

    /**
     * Fetches prefixes for the given autonomous systems and stores them in the cache.
     */
    private void fetchPrefixes(List<Integer> systems) {
        var service = context.getServiceManager().getService(BGPInfoService.class);
        systems.forEach(as -> {
            var currentProvider = service.getProvider();
            var done = false;
            do {
                SystemLogger.message("Fetching AS" + as + " prefixes...", CTX);
                var future = service.getPrefixes(as);

                while (!future.isDone() && !future.isCancelled()) {
                    Thread.yield();
                }

                try {
                    var iPv4Subnets = future.get();
                    prefixCache.put(as, iPv4Subnets);
                    done = true;
                    break;
                } catch (InterruptedException | ExecutionException e) {
                    service.fallbackNextProvider();
                    SystemLogger.error("Error happened while fetching AS" + as + " prefixes. Trying to use fallback BGP info provider:" + service.getProvider().getClass().getSimpleName(), CTX, e);
                }
            } while (service.getProvider() != currentProvider);
            if (!done) {
                SystemLogger.error("Unable to fetch AS" + as + " prefixes from all providers. Trying to use cache...", CTX);
            }
        });
    }

    /**
     * Removes event subscriptions and shuts down the executor.
     */
    @Override
    public void close() throws IOException {
        context.getEventsHandler().getSubscriptionsUpdateEvent().remove(resourceUpdateSubscription);
        context.getEventsHandler().getConfigChangeEvent().remove(configChangeSubscription);

        executor.shutdown();
    }
}
