/**
 * SBV ISO 20022 Gateway — Kafka / MQ Message Publisher
 * Phân phối message đã xác thực tới downstream consumers qua Kafka.
 *
 * Topics:
 *   sbv.iso20022.pain       — Payment Initiation (pain.*)
 *   sbv.iso20022.pacs       — Payment Clearing & Settlement (pacs.*)
 *   sbv.iso20022.camt       — Cash Management (camt.*)
 *   sbv.iso20022.auth       — Regulatory Reporting (auth.*)
 *   sbv.iso20022.mbridge    — ASEAN mBridge (pacs.009 outbound)
 *   sbv.iso20022.dlq        — Dead Letter Queue (failed messages)
 *
 * Cơ quan chủ quản: Ngân hàng Nhà nước Việt Nam (SBV)
 */

import { Kafka, Producer, RecordMetadata, logLevel } from 'kafkajs';

// ─── Config ───────────────────────────────────────────────────────────────────

const KAFKA_BROKERS = (process.env.KAFKA_BROKERS ?? 'localhost:9092').split(',');
const KAFKA_CLIENT_ID = 'sbv-iso20022-gateway';
const KAFKA_ENABLED = process.env.KAFKA_ENABLED === 'true';

/** Mapping message catalogue prefix → Kafka topic */
const TOPIC_MAP: Record<string, string> = {
  'pain': 'sbv.iso20022.pain',
  'pacs': 'sbv.iso20022.pacs',
  'camt': 'sbv.iso20022.camt',
  'auth': 'sbv.iso20022.auth',
  'head': 'sbv.iso20022.pacs',   // head.001 wraps pacs/pain
};

const MBRIDGE_TOPIC = 'sbv.iso20022.mbridge';
const DLQ_TOPIC     = 'sbv.iso20022.dlq';

// ─── Types ───────────────────────────────────────────────────────────────────

export interface PublishResult {
  published:  boolean;
  topic:      string;
  partition?: number;
  offset?:    string;
  error?:     string;
}

export interface MessageEnvelope {
  msgDefId:    string;
  uetr?:       string | null;
  senderBic?:  string;
  msgHash:     string;
  downstream:  string;
  payload:     unknown;
  receivedAt:  string;
}

// ─── Kafka producer (lazy init) ───────────────────────────────────────────────

let producer: Producer | null = null;

async function getProducer(): Promise<Producer> {
  if (producer) return producer;

  const kafka = new Kafka({
    clientId: KAFKA_CLIENT_ID,
    brokers:  KAFKA_BROKERS,
    logLevel: logLevel.WARN,
  });

  producer = kafka.producer({
    allowAutoTopicCreation: true,
    transactionTimeout: 30_000,
  });

  await producer.connect();
  return producer;
}

// ─── Public API ───────────────────────────────────────────────────────────────

/**
 * Gửi message đã xác thực vào đúng Kafka topic dựa trên msgDefId.
 * Fire-and-forget safe: nếu Kafka chưa bật, log warning và trả về published=false.
 *
 * @param envelope - Message envelope kèm metadata
 */
export async function publishMessage(envelope: MessageEnvelope): Promise<PublishResult> {
  const topic = resolveTopic(envelope.msgDefId, envelope.downstream);

  if (!KAFKA_ENABLED) {
    console.log(JSON.stringify({
      event:   'KAFKA_PUBLISH_SKIPPED',
      topic,
      msgDefId: envelope.msgDefId,
      uetr:    envelope.uetr,
      note:    'KAFKA_ENABLED=false — set env var to enable Kafka routing',
    }));
    return { published: false, topic };
  }

  try {
    const prod = await getProducer();
    const records: RecordMetadata[] = await prod.send({
      topic,
      messages: [
        {
          key:   envelope.uetr ?? envelope.msgHash,
          value: JSON.stringify(envelope),
          headers: {
            'x-sbv-msgdefid':   envelope.msgDefId,
            'x-sbv-sender-bic': envelope.senderBic ?? '',
            'x-sbv-msg-hash':   envelope.msgHash,
          },
        },
      ],
    });

    const meta = records[0];
    console.log(JSON.stringify({
      event:     'KAFKA_PUBLISHED',
      topic,
      partition: meta.partition,
      offset:    meta.offset,
      uetr:      envelope.uetr,
    }));

    return { published: true, topic, partition: meta.partition, offset: meta.offset?.toString() };

  } catch (err) {
    const error = (err as Error).message;
    console.error(JSON.stringify({ event: 'KAFKA_PUBLISH_FAILED', topic, error }));

    // Attempt DLQ
    try {
      const prod = await getProducer();
      await prod.send({
        topic: DLQ_TOPIC,
        messages: [{ key: envelope.msgHash, value: JSON.stringify({ ...envelope, dlqReason: error }) }],
      });
    } catch { /* DLQ failure — log only */ }

    return { published: false, topic, error };
  }
}

/**
 * Graceful disconnect — gọi khi server shutdown.
 */
export async function disconnectKafka(): Promise<void> {
  if (producer) {
    await producer.disconnect();
    producer = null;
  }
}

// ─── Helpers ─────────────────────────────────────────────────────────────────

function resolveTopic(msgDefId: string, downstream: string): string {
  // pacs.009 đi mBridge → topic riêng
  if (msgDefId.startsWith('pacs.009') || downstream === 'ASEAN_MBRIDGE') {
    return MBRIDGE_TOPIC;
  }
  const prefix = msgDefId.split('.')[0];
  return TOPIC_MAP[prefix] ?? 'sbv.iso20022.pacs';
}
