import { spawn } from "node:child_process";
import { randomUUID } from "node:crypto";
import path from "node:path";
import { fileURLToPath } from "node:url";
import { loadConfig } from "./config.ts";
import { PanelClient } from "./api/PanelClient.ts";
import { SessionVault } from "./session/SessionVault.ts";
import { OfficialKickApiAdapter } from "./kick/adapters/OfficialKickApiAdapter.ts";
import { parseInboundPayload, toChatMessage } from "./kick/KickWebhookEvent.ts";
import { OutboundQueue } from "./queue/OutboundQueue.ts";
import { LegacyGameCoordinator } from "./games/LegacyGameCoordinator.ts";
import type { ChannelStatus, ChatMessage, EngineEvent, InboundKickEvent } from "./types.ts";

const config = loadConfig();
const panel = new PanelClient(config);
const vault = new SessionVault(config.vaultPath, config.oauthMasterKey);
const adapter = new OfficialKickApiAdapter(config, vault);
const channels = new Map<number, ChannelStatus>();
let stopping = false;
let adapterReady = false;
let botKickUserId = config.botKickUserId;
let channelsSnapshot = 0;
// Commands must never queue behind routine chat activity (points, chat
// overlay, emotes…).  They use a dedicated small lane while the ordinary
// activity lane can batch safely in the background.
let priorityEventBuffer: EngineEvent[] = [];
let activityEventBuffer: EngineEvent[] = [];
let activityEventFlush: NodeJS.Timeout | null = null;
let priorityFlushInFlight: Promise<void> | null = null;
let activityFlushInFlight: Promise<void> | null = null;
let activityAcknowledgementInFlight: Promise<void> = Promise.resolve();
let lastKickSubscriptionMaintenanceAt = 0;
let lastRuntimeConfigSyncAt = 0;

const flushEvents = async (strict = false, priorityOnly = false): Promise<void> => {
  const buffer = priorityOnly ? priorityEventBuffer : activityEventBuffer;
  const inFlight = priorityOnly ? priorityFlushInFlight : activityFlushInFlight;
  // Two command flushes keep their order. Activity is deliberately a
  // different lane: it must not make an incoming !command wait on a batch of
  // viewer-point or overlay updates.
  if (inFlight) {
    try {
      await inFlight;
    } catch (error) {
      if (strict) throw error;
      return;
    }
  }
  if (!buffer.length) return;
  const batch = buffer.splice(0, priorityOnly ? 12 : 100);
  const publish = async (): Promise<void> => {
    try {
      await panel.events(batch);
    } catch (error) {
      buffer.unshift(...batch);
      // An unavailable panel must not turn into a huge replay backlog that
      // answers old chat commands after recovery. The webhook inbox remains
      // the durable delivery source; this short buffer only bridges a brief
      // transient HTTP failure.
      if (buffer.length > 300) {
        if (priorityOnly) priorityEventBuffer = buffer.slice(-300);
        else activityEventBuffer = buffer.slice(-300);
      }
      if (strict) throw error;
    }
  };
  const request = publish();
  if (priorityOnly) priorityFlushInFlight = request;
  else activityFlushInFlight = request;
  try {
    await request;
  } finally {
    if (priorityOnly && priorityFlushInFlight === request) priorityFlushInFlight = null;
    if (!priorityOnly && activityFlushInFlight === request) activityFlushInFlight = null;
  }
};

const enqueueEvent = (event: EngineEvent, priority = false): void => {
  if (priority) {
    priorityEventBuffer.push(event);
    return;
  }
  activityEventBuffer.push(event);
  if (activityEventFlush) return;
  activityEventFlush = setTimeout(() => {
    activityEventFlush = null;
    void flushEvents();
  }, 80);
  activityEventFlush.unref();
};

// Activity delivery is durable through kick_webhook_inbox.  Keep its final
// panel write + acknowledgement in order, but off the command loop: viewer
// chat activity must never delay the next !command received by Kick.
const scheduleActivityAcknowledgements = (events: InboundKickEvent[]): void => {
  if (!events.length) return;
  const deliver = async (): Promise<void> => {
    try {
      await flushEvents(true);
      for (const event of events) {
        await panel.acknowledgeInboundEvent(event.id, "processed");
      }
    } catch (error) {
      const detail = error instanceof Error ? error.message : "Publication des événements indisponible";
      for (const event of events) {
        await panel.acknowledgeInboundEvent(event.id, event.attempts >= 5 ? "failed" : "retry", detail).catch(() => undefined);
      }
    }
  };
  activityAcknowledgementInFlight = activityAcknowledgementInFlight
    .catch(() => undefined)
    .then(deliver);
};

const engineInstanceId = randomUUID();
const games = new LegacyGameCoordinator(
  config,
  enqueueEvent,
  (channelId, gameSessionId, action) => panel.gameLease(channelId, gameSessionId, engineInstanceId, action)
);

const handleMessage = (message: ChatMessage): void => {
  if (
    (botKickUserId && message.sender.kick_user_id === botKickUserId)
    || message.sender.username.toLocaleLowerCase("fr") === config.botUsername.toLocaleLowerCase("fr")
  ) return;
  enqueueEvent({
    id: message.event_id,
    channel_id: message.channel_id,
    type: "chat.message",
    payload: {
      content: message.content,
      sender: message.sender,
      kick_channel_id: message.kick_channel_id,
      kick_chatroom_id: message.kick_chatroom_id,
      created_at: message.created_at
    }
  }, /^![a-z0-9_-]{1,40}(?:\s|$)/i.test(message.content.trim()));
  games.handleMessage(message);
};

const queue = new OutboundQueue(panel, adapter, (channelId) => channels.get(channelId));

// Queue polling must not be chained behind the inbound-event request. The
// latter is a PHP/API call and, during a transient timeout, used to hold every
// pending chat reply for up to five seconds. This independent poller keeps the
// outbound Kick path moving while the maintenance path heals.
const kickQueueNow = (): void => {
  if (!adapterReady || stopping) return;
  void queue.poll().catch(() => undefined);
};

const initializeSession = async (): Promise<void> => {
  try {
    await adapter.initialize();
    const session = await adapter.testSession();
    adapterReady = session.valid;
    if (session.valid && session.kickUserId) botKickUserId = session.kickUserId;
    if (!adapterReady) queue.openCircuit(120_000);
    // The panel status badge is observability only.  A transient response
    // problem must never disable a valid Kick session or freeze every chat
    // reply in the outbound queue.
    await panel.sessionStatus(session.valid ? "valid" : "invalid", {
      kick_user_id: session.kickUserId,
      expires_at: session.expiresAt,
      error: session.error ?? ""
    }).catch((error) => {
      const detail = error instanceof Error ? error.message : "Statut de session indisponible";
      process.stderr.write(`[FilyX] statut de session ignoré : ${detail}\n`);
    });
  } catch (error) {
    adapterReady = false;
    queue.openCircuit(120_000);
    const detail = error instanceof Error ? error.message : "Autorisation OAuth absente";
    await panel.sessionStatus("invalid", { error: detail }).catch(() => undefined);
  }
};

const syncChannels = async (signal: AbortSignal): Promise<void> => {
  const nextChannels = await panel.channels(signal);
  // Keep the current map available while the panel request is in flight. An
  // empty map here made freshly received Kick webhooks retry needlessly.
  const nextChannelMap = new Map<number, ChannelStatus>();
  for (const channel of nextChannels) nextChannelMap.set(channel.channel_id, channel);
  channels.clear();
  for (const [channelId, channel] of nextChannelMap) channels.set(channelId, channel);
  channelsSnapshot = nextChannels.length;
  const now = Date.now();
  // Webhook subscription reconciliation talks to Kick for every linked
  // channel. Running it every ten seconds made the engine spend more time on
  // background maintenance than on chat, and amplified any transient Kick
  // latency. Keep a safe two-minute refresh interval instead; new game and
  // command activity still synchronises immediately on its own path.
  if (adapterReady && now - lastKickSubscriptionMaintenanceAt >= 120_000) {
    lastKickSubscriptionMaintenanceAt = now;
    // Kick can temporarily rate-limit subscription reconciliation. Incoming
    // webhooks and queued chat replies must keep running even when this
    // maintenance step fails, otherwise one HTTP 429 freezes every channel.
    await adapter.ensureBotChannelSubscription(Number(botKickUserId)).catch((error) => {
      const message = error instanceof Error ? error.message : "abonnement du bot indisponible";
      process.stderr.write(`[FilyX] abonnement bot reporté : ${message}\n`);
    });
    await adapter.ensureEventSubscriptions(nextChannels).catch((error) => {
      const message = error instanceof Error ? error.message : "abonnements de chaînes indisponibles";
      process.stderr.write(`[FilyX] abonnements chaînes reportés : ${message}\n`);
    });
  }
  const refreshRuntimeConfig = now - lastRuntimeConfigSyncAt >= 20_000;
  if (refreshRuntimeConfig) lastRuntimeConfigSyncAt = now;
  for (const channel of nextChannels) {
    if (refreshRuntimeConfig) {
      const runtimeConfig = await panel.channelConfig(channel.channel_id, signal).catch((error) => {
        const message = error instanceof Error ? error.message : "configuration de chaîne indisponible";
        process.stderr.write(`[FilyX] synchronisation jeux ignorée : ${message}\n`);
        return null;
      });
      // A temporary panel error must never be interpreted as "aucun jeu".
      // Otherwise it destroys the live engine and a refresh starts it again.
      if (runtimeConfig) {
        await games.sync(channel, (runtimeConfig.active_game as Parameters<typeof games.sync>[1]) ?? null);
      }
    }
    if (refreshRuntimeConfig && channel.moderator_status === "checking" && adapterReady) {
      const status = await adapter.checkModeratorRole(channel);
      await panel.moderatorStatus(channel.channel_id, status, botKickUserId);
    }
  }
};

const processInboundEvent = async (event: InboundKickEvent): Promise<void> => {
  const channel = channels.get(event.channel_id);
  if (!channel) throw new Error(`Chaîne ${event.channel_id} non chargée.`);
  const payload = parseInboundPayload(event);
  const chatMessage = toChatMessage(event, channel, payload);
  if (chatMessage) {
    handleMessage(chatMessage);
    return;
  }
  enqueueEvent({
    id: event.kick_message_id,
    channel_id: event.channel_id,
    type: event.event_type,
    payload
  });
};

const isPriorityChatCommand = (event: InboundKickEvent): boolean => {
  if (event.event_type !== "chat.message.sent") return false;
  try {
    return /^![a-z0-9_-]{1,40}(?:\s|$)/i.test(String(parseInboundPayload(event).content ?? "").trim());
  } catch {
    return false;
  }
};

const gameCommandAction = (event: InboundKickEvent): "start" | "stop" | null => {
  if (event.event_type !== "chat.message.sent") return null;
  try {
    const content = String(parseInboundPayload(event).content ?? "").trim();
    if (/^!stop$/i.test(content)) return "stop";
    return /^!(?:quiz|motus|quisuisje|whoami|petitbac)(?:\s|$)/i.test(content) ? "start" : null;
  } catch {
    return null;
  }
};

// Game sessions are written by the panel while the inbound event is flushed.
// Fetching just that channel after the confirmed write removes the former
// 10-second wait for the maintenance synchroniser without starting a second
// engine: LegacyGameCoordinator still owns both the session and DB lease.
const syncGameForChannel = async (channelId: number): Promise<void> => {
  const channel = channels.get(channelId);
  if (!channel) return;
  const runtimeConfig = await panel.channelConfig(channelId);
  await games.sync(channel, (runtimeConfig.active_game as Parameters<typeof games.sync>[1]) ?? null);
};

const pollInboundEvents = async (signal: AbortSignal): Promise<void> => {
  const events = await panel.inboundEvents(signal);
  // A busy chat can contain dozens of ordinary messages. Commands must not
  // wait for ordinary activity.  The command path is deliberately flushed
  // one event at a time; the remaining chat messages are published in a
  // single panel request below.  Previously every normal message paid for a
  // full HTTP/PHP request before the next inbound poll, which made commands
  // feel slow during an active chat despite their database priority.
  const orderedEvents = [...events].sort((left, right) => Number(isPriorityChatCommand(right)) - Number(isPriorityChatCommand(left)));
  const priorityEvents = orderedEvents.filter(isPriorityChatCommand);
  const activityEvents = orderedEvents.filter((event) => !isPriorityChatCommand(event));

  const processAndAcknowledge = async (event: InboundKickEvent, priorityCommand: boolean, flushImmediately: boolean): Promise<boolean> => {
    const gameAction = gameCommandAction(event);
    try {
      await processInboundEvent(event);
      if (flushImmediately) await flushEvents(true, true);
      if (gameAction === "start") {
        await syncGameForChannel(event.channel_id).catch((error) => {
          const detail = error instanceof Error ? error.message : "démarrage du mini-jeu reporté";
          process.stderr.write(`[FilyX] synchronisation immédiate jeu ignorée : ${detail}\n`);
        });
      } else if (gameAction === "stop") {
        await games.stopChannel(event.channel_id);
      }
      await panel.acknowledgeInboundEvent(event.id, "processed");
      // Request a queue drain immediately, without making the inbound
      // acknowledgement wait for a stale queue request.
      if (priorityCommand && adapterReady) kickQueueNow();
      return true;
    } catch (error) {
      const detail = error instanceof Error ? error.message : "Webhook Kick invalide";
      await panel.acknowledgeInboundEvent(event.id, event.attempts >= 5 ? "failed" : "retry", detail);
      return false;
    }
  };

  // Commands stay on the smallest possible path: persist, acknowledge and
  // request the outgoing Kick queue right away.
  for (const event of priorityEvents) {
    await processAndAcknowledge(event, true, true);
  }

  // Standard chat activity still reaches points, overlays and games, but no
  // longer starts a separate PHP request per message.  The batch is durable
  // before it is acknowledged, so a temporary panel failure remains safely
  // retryable by the webhook inbox.
  const processedActivity: InboundKickEvent[] = [];
  for (const event of activityEvents) {
    try {
      await processInboundEvent(event);
      processedActivity.push(event);
    } catch (error) {
      const detail = error instanceof Error ? error.message : "Webhook Kick invalide";
      await panel.acknowledgeInboundEvent(event.id, event.attempts >= 5 ? "failed" : "retry", detail);
    }
  }
  scheduleActivityAcknowledgements(processedActivity);
};

const processWorkerCommands = async (signal: AbortSignal): Promise<void> => {
  const commands = await panel.workerCommands(signal);
  for (const command of commands) {
    let status: "completed" | "failed" = "completed";
    let result = "";
    try {
      if (command.command_type === "test-session") {
        const session = await adapter.testSession();
        result = session.valid ? "Autorisation OAuth FiIyx valide." : (session.error ?? "Autorisation invalide.");
        if (!session.valid) status = "failed";
      } else if (command.command_type === "revoke-session") {
        await adapter.close();
        await vault.revoke();
        adapterReady = false;
        queue.openCircuit(24 * 60 * 60 * 1_000);
        result = "Autorisation OAuth chiffrée révoquée localement.";
      } else if (command.command_type === "open-session-manager") {
        if (process.env.FILYX_ALLOW_INTERACTIVE_SESSION_MANAGER !== "true") {
          throw new Error("L’assistant OAuth doit être lancé manuellement sur la machine du moteur.");
        }
        const here = path.dirname(fileURLToPath(import.meta.url));
        const child = spawn(process.execPath, ["--experimental-strip-types", path.join(here, "session", "session-manager.ts")], {
          detached: true,
          stdio: "ignore",
          windowsHide: false,
          env: process.env
        });
        child.unref();
        result = "Assistant OAuth lancé.";
      } else if (command.command_type === "restart-engine") {
        result = "Redémarrage demandé.";
        await panel.completeWorkerCommand(command.id, status, result);
        setTimeout(() => process.exit(0), 200).unref();
        continue;
      } else {
        throw new Error("Commande worker non prise en charge.");
      }
    } catch (error) {
      status = "failed";
      result = error instanceof Error ? error.message : "Erreur inconnue";
    }
    await panel.completeWorkerCommand(command.id, status, result);
  }
};

const loop = async (): Promise<void> => {
  const controller = new AbortController();
  const heartbeat = setInterval(() => {
    void panel.heartbeat({
      active_channels: channelsSnapshot,
      websocket_connections: 0,
      queue_depth: priorityEventBuffer.length + activityEventBuffer.length,
      metadata: {
        adapter_ready: adapterReady,
        queue_circuit_open: queue.circuitOpen,
        kick_transport: "official-api-webhooks"
      }
    }).catch(() => undefined);
  }, config.heartbeatIntervalMs);
  heartbeat.unref();
  // Inbound commands explicitly trigger an immediate drain. When idle, a
  // 300-ms poll is already imperceptible to chat and avoids needlessly
  // hammering PHP/MariaDB eight times a second.
  const outboundPoller = setInterval(kickQueueNow, Math.max(250, Math.min(400, config.pollIntervalMs)));
  outboundPoller.unref();
  // Streamlabs remains the payment provider. The engine only asks the panel
  // for newly verified Streamlabs donation events, then the usual alert/chat
  // pipeline handles them. Keeping it apart from the Kick command loop means
  // a slow third-party call can never delay a chat response.
  let streamlabsSyncInFlight = false;
  const streamlabsPoller = setInterval(() => {
    if (stopping || streamlabsSyncInFlight) return;
    streamlabsSyncInFlight = true;
    void panel.syncStreamlabsDonations()
      .catch(() => undefined)
      .finally(() => { streamlabsSyncInFlight = false; });
  }, 5_000);
  streamlabsPoller.unref();

  let lastChannelSync = 0;
  let lastCommandPoll = 0;
  let channelSyncInFlight: Promise<void> | null = null;
  let workerCommandPollInFlight = false;
  const scheduleChannelSync = (): void => {
    if (channelSyncInFlight) return;
    lastChannelSync = Date.now();
    channelSyncInFlight = syncChannels(controller.signal)
      .catch((error) => {
        const message = error instanceof Error ? error.message : "Synchronisation des chaînes impossible";
        process.stderr.write(`[FilyX] ${message}\n`);
      })
      .finally(() => { channelSyncInFlight = null; });
  };
  // These commands operate the engine itself (OAuth test, restart…) and do
  // not affect an incoming Kick message.  They used to be awaited directly
  // in the same loop as inbound webhooks: one slow PHP request could then
  // hold chat processing for up to five seconds.  Keep one background poll
  // in flight, exactly as the other maintenance tasks do.
  const scheduleWorkerCommandPoll = (): void => {
    if (workerCommandPollInFlight) return;
    workerCommandPollInFlight = true;
    void processWorkerCommands(controller.signal)
      .catch((error) => {
        const message = error instanceof Error ? error.message : "Commandes moteur indisponibles";
        process.stderr.write(`[FilyX] ${message}\n`);
      })
      .finally(() => { workerCommandPollInFlight = false; });
  };
  // Load the first channel map before polling Kick. Later maintenance runs in
  // the background, so a slow subscription reconciliation never freezes chat.
  scheduleChannelSync();
  await channelSyncInFlight;
  while (!stopping) {
    const started = Date.now();
    try {
      if (started - lastChannelSync > 10_000) scheduleChannelSync();
      await pollInboundEvents(controller.signal);
      // Outbound messages have their own low-latency poller above. Do not
      // await them here: a delayed inbound poll must not freeze chat replies.
      kickQueueNow();
      if (started - lastCommandPoll > 1_000) {
        lastCommandPoll = started;
        scheduleWorkerCommandPoll();
      }
      // Routine messages are flushed independently. Do not let an activity
      // batch (points, overlays, emotes…) pause the next inbound command
      // poll; priority commands flush in their own lane above.
      void flushEvents();
    } catch (error) {
      const message = error instanceof Error ? error.message : "Erreur moteur";
      process.stderr.write(`[FilyX] ${message}\n`);
    }
    const elapsed = Date.now() - started;
    await new Promise((resolve) => setTimeout(resolve, Math.max(100, config.pollIntervalMs - elapsed)));
  }
  controller.abort();
  clearInterval(heartbeat);
  clearInterval(outboundPoller);
  clearInterval(streamlabsPoller);
};

const shutdown = async (): Promise<void> => {
  if (stopping) return;
  stopping = true;
  if (activityEventFlush) clearTimeout(activityEventFlush);
  await flushEvents(true, true).catch(() => undefined);
  await flushEvents(true).catch(() => undefined);
  channels.clear();
  games.stop();
  await adapter.close();
};

process.once("SIGINT", () => void shutdown());
process.once("SIGTERM", () => void shutdown());

await initializeSession();
await loop();
