<?php
namespace App\Core\Service;
use App\Config\Main;
use RdKafka\Producer;
use RdKafka\KafkaConsumer;
use RdKafka\Conf;
require 'vendor/autoload.php';

class Kafka
{
    public static $bootstrap_servers = "localhost:9092";
    public static $topic_name = "sansoftwares";
    public static $producer    = NULL;
    public static $consumer    = NULL;

    private function log($message) {
        try {
            if($_ENV['LOG_PATH']) {
                if(file_exists($_ENV['LOG_PATH']. '/' . date('Y-m-d') . '.log')) {
                    error_log("\r\n".$message, 3, $_ENV['LOG_PATH']. '/' . date('Y-m-d') . '.log');
                } else {
                    error_log($message);
                }
            } else {
                error_log($message);
            }
        } catch(\Exception $e) {
            error_log($e->getMessage());
        }
    }

    private static function setInitVars() {
        try {
            $bs = Main::getKafkaBootstrapServer();
            if(!empty($bs)) {
                self::$bootstrap_servers = $bs;
            }
            $topic = Main::getKafkaTopic();
            if(!empty($bs)) {
                self::$topic_name = $topic;
            }
            return true;
        } catch(\Exception $e) {
            self::log($e->getMessage());
            return false;
        }
    }

    public static function setProducer()
    {
        try {
            if(!self::$producer) {
                $initialized = self::setInitVars();
                if(empty($initialized)) {
                    throw new \Exception("Unable to set Kafka configuration");
                }
                self::$producer = new Producer();
                self::$producer->addBrokers(self::$bootstrap_servers);
                self::$producer->newTopic(self::$topic_name);
            }
        } catch(\Exception $e) {
            self::log($e->getMessage());
        }
    }

    public static function setConsumer()
    {
        try {
            if(!self::$consumer) {
                $initialized = self::setInitVars();
                if(empty($initialized)) {
                    throw new \Exception("Unable to set Kafka configuration");
                }
                $conf = new Conf();
                $conf->set('group.id', 'server-consumer-group');
                $conf->set('auto.offset.reset', 'earliest');
                self::$consumer = new KafkaConsumer($conf);
                // self::$consumer->subscribe(self::$topic_name);
            }
        } catch(\Exception $e) {
            self::log($e->getMessage());
        }
    }
}
