<?php
namespace App\Service;
use App\Core\Service\Logger;
use App\Config\Main;
use RdKafka;
class Kafkam
{
    public static $config;
    public static $producer;
    public static $admin;
    public static $topics = [];

    public function __construct()
    {
        if(empty(self::$config)) {
            self::initialize();           
        }
    }

    public static function initialize() {
        try {
            $broker = Main::getKafkaBroker();
            if(empty($broker)) {
                throw new \Exception("Kafka Broker not configured");
            }
            self::$config   = new RdKafka\Conf();
            self::$config->set('bootstrap.servers', $broker);
            $topic_list   = Main::KafkaTopics();
            if(!empty($topic_list)) {
                foreach($topic_list as $topic) {
                    self::$topics[$topic] = ["name"=>$topic, "exists"=>false, "instance"=>NULL];
                }
            }
            self::$producer = new RdKafka\Producer(self::$config);
            $metadata = self::$producer->getMetadata(true, null, 3000);
            $kafka_topics = $metadata->getTopics();

            foreach ($kafka_topics as $t) {
                $topic_name = $t->getTopic();
                if(array_key_exists($topic_name, self::$topics)) {
                    self::$topics[$topic_name]['exists'] = true;
                    self::$topics[$topic_name]['instance'] = $t;
                }
            }
            foreach (self::$topics as $topic_name=>$topic) {
                if($topic['exists'] == false) {
                    $topicConf = new RdKafka\TopicConf();
                    $new_topic = self::$producer->newTopic($topic_name, $topicConf);
                    self::$topics[$topic_name]['instance']  = $new_topic;
                }
            }
        } catch(\Exception $e) {
            Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
            return ["status"=>false, "error"=>$e->getMessage()];
        }
    }

    public function sendMessage(String $topic_name, Array $data) {
        try {
            if(empty(self::$producer)) {
                throw new \Exception("Kafka Producer not configured");
            }
            if(!array_key_exists($topic_name, self::$topics)) {
                $topicConf = new RdKafka\TopicConf();
                $new_topic = self::$producer->newTopic($topic_name, $topicConf);
                self::$topics[$topic_name]['name'] = $topic_name;
                self::$topics[$topic_name]['exists'] = true;
                self::$topics[$topic_name]['instance']  = $new_topic;
            }
            $new_topic = self::$producer->newTopic($topic_name);
            $new_topic->produce(RD_KAFKA_PARTITION_UA, 0, json_encode($data));
            self::$producer->flush(10000);
            return ["status"=>true];
        } catch(\Exception $e) {
            Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
            return ["status"=>false, "error"=>$e->getMessage()];
        }
    }
} 