import type { KickSessionAdapter } from "../kick/adapters/KickSessionAdapter.ts";
import type { PanelClient } from "../api/PanelClient.ts";
import type { ChannelStatus, OutboundMessage } from "../types.ts";

export class OutboundQueue {
  private readonly panel: PanelClient;
  private readonly adapter: KickSessionAdapter;
  private readonly channels: (channelId: number) => ChannelStatus | undefined;
  private readonly delivered = new Map<string, number>();
  private readonly channelSentAt = new Map<number, number>();
  private globalSentAt = 0;
  private processing = false;
  private pollRequested = false;
  private circuitOpenUntil = 0;

  constructor(panel: PanelClient, adapter: KickSessionAdapter, channels: (channelId: number) => ChannelStatus | undefined) {
    this.panel = panel;
    this.adapter = adapter;
    this.channels = channels;
  }

  async poll(signal?: AbortSignal): Promise<void> {
    if (Date.now() < this.circuitOpenUntil) return;
    if (this.processing) {
      // The engine polls every few hundred milliseconds. Keep exactly one
      // follow-up poll when a message is currently in the Kick rate limiter,
      // rather than losing a newly queued viewer command until a later loop.
      this.pollRequested = true;
      return;
    }
    this.processing = true;
    try {
      const messages = await this.panel.queue(signal);
      for (const message of messages) {
        if (signal?.aborted) break;
        await this.process(message);
      }
      this.pruneDelivered();
    } finally {
      this.processing = false;
      if (this.pollRequested && Date.now() >= this.circuitOpenUntil) {
        this.pollRequested = false;
        queueMicrotask(() => { void this.poll().catch(() => undefined); });
      }
    }
  }

  openCircuit(milliseconds = 60_000): void {
    this.circuitOpenUntil = Math.max(this.circuitOpenUntil, Date.now() + milliseconds);
  }

  get circuitOpen(): boolean {
    return Date.now() < this.circuitOpenUntil;
  }

  private async process(message: OutboundMessage): Promise<void> {
    if (this.delivered.has(message.idempotency_key)) {
      await this.panel.acknowledge(message, "sent");
      return;
    }
    const channel = this.channels(message.channel_id);
    if (!channel) {
      await this.panel.acknowledge(message, "retry", "Chaîne non chargée par le worker.");
      return;
    }
    await this.waitForRateLimit(channel.channel_id);
    try {
      await this.adapter.sendChatMessage(channel, message.message_text, message.reply_to_message_id);
      const now = Date.now();
      this.globalSentAt = now;
      this.channelSentAt.set(channel.channel_id, now);
      this.delivered.set(message.idempotency_key, now);
      await this.panel.acknowledge(message, "sent");
    } catch (error) {
      const detail = error instanceof Error ? error.message : "Erreur d’envoi inconnue";
      if (/session|login|auth/i.test(detail)) this.openCircuit(120_000);
      await this.panel.acknowledge(message, message.attempts >= 5 ? "failed" : "retry", detail);
    }
  }

  private async waitForRateLimit(channelId: number): Promise<void> {
    const now = Date.now();
    const globalWait = Math.max(0, 650 - (now - this.globalSentAt));
    const channelWait = Math.max(0, 1_250 - (now - (this.channelSentAt.get(channelId) ?? 0)));
    const delay = Math.max(globalWait, channelWait);
    if (delay > 0) await new Promise((resolve) => setTimeout(resolve, delay));
  }

  private pruneDelivered(): void {
    const threshold = Date.now() - 3_600_000;
    for (const [key, timestamp] of this.delivered) {
      if (timestamp < threshold) this.delivered.delete(key);
    }
  }
}
