<?php
namespace App\Worker;

ini_set('display_errors', 1);
error_reporting(E_ALL);

use App\Config\Main;
use RdKafka;
require_once dirname(__DIR__, 2) . '/vendor/autoload.php';
$dotenv = \Dotenv\Dotenv::createImmutable(dirname(__DIR__, 2));
$dotenv->load();
require_once dirname(__DIR__, 2) . '/app/Core/Service/Logger.php';
require_once dirname(__DIR__, 2) . '/app/Helpers/common.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/OrderController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/PinController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/DocketController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/PaymentController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/ShopifyController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/MailerController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Client/LeadController.php';
require_once dirname(__DIR__, 2). '/app/Controllers/Admin/NotificationController.php';
use App\Controllers\Admin\NotificationController;
use App\Controllers\Client\PinController;
use App\Controllers\Client\OrderController;
use App\Controllers\Client\DocketController;
use App\Controllers\Client\ShopifyController;
use App\Controllers\Client\PaymentController;
use App\Controllers\Client\MailerController;
use App\Controllers\Client\LeadController;
use App\Core\Service\Logger;

function consumeSentEmailSMS($argv) {
    try {
        $instance = new MailerController();
        $resp = $instance->processEmailSMSLog($argv['comp_id'], $argv['user_id'], $argv['req_id']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeUpdateOrderStatus($argv) {
    try {
        $instance = new OrderController();
        $resp = $instance->processUpdateOrderStatus($argv['comp_id'], $argv['req_id']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeSavePincodes($argv) {
    try {
        $instance = new PinController();
        $resp = $instance->savePincode($argv['comp_id'], $argv['user_id'], $argv['courier_id'], $argv['file_path']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeSaveDocketNo($argv) {
    try {
        $instance = new DocketController();
        $resp = $instance->saveDocketNo($argv['comp_id'], $argv['user_id'], $argv['courier_id'], $argv['file_path']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function counsumeSendNotification($argv) {
    try {
        $instance = new NotificationController();
        $resp = $instance->sendPendingNotification($argv);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeSendPaymentLink($argv) {
    try {
        $instance = new PaymentController();
        $resp = $instance->processReqToSendPaymentLink($argv['comp_id'], $argv['req_id']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeMoveOrderToShopify($argv) {
    try {
        $instance = new ShopifyController();
        $resp = $instance->moveOrderToShopify($argv['comp_id']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeSendNotification($argv) {
    try {
        $instance = new NotificationController();
        $resp = $instance->sendPendingNotification($argv['log_id']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

function consumeCreateOrderFromLead($argv) {
    try {
        $instance = new LeadController();
        $resp = $instance->createOrderFromLead($argv['comp_id'], $argv['user_id'], $argv['req_id']);
        if(empty($resp['status'])) {
            throw new \Exception($resp['error']);
        }
        return true;
    } catch(\Exception $e){
        Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        return false;
    }
}

try {
    $broker = Main::getKafkaBroker();

    if(empty($broker)) {
        throw new \Exception("Kafka Broker not configured in consumer");
    }
    $conf = new RdKafka\Conf();
    $conf->set('bootstrap.servers', $broker);  
    $conf->set('group.id', 'local_sanordo_consumer_group');    
    $conf->set('auto.offset.reset', 'earliest');
    $conf->set('enable.auto.commit', 'false');
    $conf->set('max.poll.interval.ms', '600000');
    $conf->set('session.timeout.ms', '30000');
    $conf->set('heartbeat.interval.ms', '3000');
    // $conf->set('log_level', (string) LOG_DEBUG);
    // $conf->set('debug', 'all');
    $consumer = new RdKafka\KafkaConsumer($conf);
    $consumer->subscribe(Main::KafkaTopics());
    
    echo "\nWaiting for messages...\n";
    
    while (true) {
        try {
            $message = $consumer->consume(1000); // Poll for messages with a 120-second timeout
            if ($message === null) {
                continue;
            }

            switch ($message->err) {
                case \RD_KAFKA_RESP_ERR_NO_ERROR:
                    // echo sprintf("Message received: Topic '%s', Partition %d, Offset %d, Payload: %s\n",
                    // $message->topic_name, $message->partition, $message->offset, $message->payload);
                    // echo $message->topic_name;
                    // print_r($message->payload);
                    Logger::debug($message->topic_name, ['File' => __FILE__, 'Line'=>__LINE__]);
                    if(in_array($message->topic_name, Main::KafkaTopics())){
                        $args = json_decode($message->payload, true);
                        if (json_last_error() !== JSON_ERROR_NONE) {
                            throw new \Exception("Invalid JSON payload");
                        }
                        $processed = false;
                        if($message->topic_name == 'sanordo_update_order_status') {
                            $processed = consumeUpdateOrderStatus($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeUpdateOrderStatus");
                            }
                        }
                        if($message->topic_name == 'sanordo_save_servicable_pincodes') {
                            $processed = consumeSavePincodes($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeSavePincodes");
                            }
                        }
                        if($message->topic_name == 'sanordo_save_docketno') {
                            $processed = consumeSaveDocketNo($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeSaveDocketNo");
                            }
                        }
                        if($message->topic_name == 'sanordo_share_payment_link') {
                            $processed = consumeSendPaymentLink($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeSendPaymentLink");
                            }
                        }
                        if($message->topic_name == 'sanordo_move_order_shopify') {
                            $processed = consumeMoveOrderToShopify($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeMoveOrderToShopify");
                            }
                        }
                        if($message->topic_name == 'sanordo_send_admin_notifications') {
                            $processed = consumeSendNotification($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeSendNotification");
                            }
                        }
                        if($message->topic_name == 'sanordo_client_email_sms_whatsapp') {
                            $processed = consumeSentEmailSMS($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeSentEmailWhatsApp");
                            }
                        }
                        if($message->topic_name == 'sanordo_create_order_from_lead') {
                            $processed = consumeCreateOrderFromLead($args);
                            if(!$processed){
                                throw new \Exception("Could not consumeCreateOrderFromLead");
                            }
                        }
                        if($processed) {
                            $consumer->commit($message);
                        }
                    }
                    break;
                case \RD_KAFKA_RESP_ERR__PARTITION_EOF:
                    echo "No more messages for this partition, waiting for new messages...\n";
                    break;
                case \RD_KAFKA_RESP_ERR__TIMED_OUT:
                    // echo "Consumer timed out, no messages received within the timeout period.\n";
                    break;
                default:
                    throw new \Exception($message->errstr(), $message->err);
            }
        } catch(\Exception $e) {
            Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
        }
    }
} catch(\Exception $e) {
    Logger::error($e->getMessage(), ['File' => $e->getFile(), 'Line'=>$e->getLine()]);
}
?>