from confluent_kafka import Consumer, Producer
import logging
import ssl
import json
import datetime
from dateutil import parser
from datetime import datetime, timedelta
import pandas as pd
import time
import json
import re
from configparser import ConfigParser
from elasticsearch import Elasticsearch, helpers


logging.basicConfig(filename='/var/www/html/kafka_consumer/logs/deywuro_logs.log',
                    filemode='a',
                    format='%(asctime)s %(name)s %(levelname)s %(message)s',
                    level=logging.INFO)

#Read config.ini file
config_object = ConfigParser()
config_object.read("config.ini")



#Get info
elasticinfo = config_object["elastic"]

consumer = Consumer({
  "bootstrap.servers": "static.66.8.9.5.clients.your-server.de:9092,static.79.14.243.136.clients.your-server.de:9092,static.171.45.217.95.clients.your-server.de:9092",
  'security.protocol': 'SASL_SSL',
  'sasl.mechanisms': "PLAIN",
  'sasl.username': "client1",
  'sasl.password': "Pay50ToconFrom300Am200aM",
  'ssl.ca.location': "/var/www/html/kafka/certs/my_key.pem",
  'group.id': 'deywuro-logs-0820',
  'auto.offset.reset': 'earliest'
})



try:
    consumer.subscribe(['hellio_live_c1.hellio.logs'])
    new_list=[]
    now_time = datetime.now()
    es = Elasticsearch([str(elasticinfo['local'])])
    while True:
        try:
          if datetime.now() > now_time + timedelta(seconds=90):
            res = helpers.bulk(es, new_list)#,raise_on_error=False)
            logging.info(res)
            now_time= datetime.now()
            new_list = []
          else:
            pass 
        except Exception as e:
          logging.info('%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%')
          logging.info(e) 
          logging.info('%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%')
          continue
        msg = consumer.poll(timeout=1.0)
        if not msg: 
          continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                # End of partition event
                sys.stderr.write('%% %s [%d] reached end at offset %d\n' %
                                 (msg.topic(), msg.partition(), msg.offset()))
            elif msg.error():
                raise KafkaException(msg.error())
        else:
            try:
                pay_load = json.loads(msg.value().decode('utf-8'))
                if not pay_load['after']['created_at'] or pay_load['after']['created']=='None':
                  created_at = datetime.utcnow().date()
                  created_at = created_at.strftime('%Y-%m-%d')
                  created_at = datetime.strptime(created_at, '%Y-%m-%d')
                else:
                  created_at = parser.parse(pay_load['after']['created_at']).isoformat()
                  created_at = created_at.replace('T', ' ')
                  created_at = created_at.replace(created_at[19:], '')
                  created_at = datetime.strptime(created_at, '%Y-%m-%d %H:%M:%S')

                if not pay_load['after']['updated_at']:
                  updated_at = datetime.utcnow().date()
                  updated_at = updated_at.strftime('%Y-%m-%d')
                  updated_at = datetime.strptime(updated_at, '%Y-%m-%d')
                else:
                  updated_at = parser.parse(pay_load['after']['updated_at']).isoformat()
                  updated_at = updated_at.replace('T', ' ')
                  updated_at = updated_at.replace(updated_at[19:], '')
                  updated_at = datetime.strptime(updated_at, '%Y-%m-%d %H:%M:%S')


                if not pay_load['after']['submit_date']:
                  submit_date= datetime.utcnow().strftime('%Y-%m-%d %H:%M:%S')
                  submit_date = datetime.strptime(submit_date, '%Y-%m-%d %H:%M:%S')
                else:
                  submit_date = datetime.utcfromtimestamp((pay_load['after']['submit_date'])/1000).strftime('%Y-%m-%d %H:%M:%S')
                  submit_date = datetime.strptime(submit_date, '%Y-%m-%d %H:%M:%S')

                if not pay_load['after']['delivery_date']:
                  delivery_date= datetime.utcnow().strftime('%Y-%m-%d %H:%M:%S')
                  delivery_date = datetime.strptime(delivery_date, '%Y-%m-%d %H:%M:%S')
                else:
                  delivery_date = datetime.utcfromtimestamp((pay_load['after']['delivery_date'])/1000).strftime('%Y-%m-%d %H:%M:%S')
                  delivery_date = datetime.strptime(delivery_date, '%Y-%m-%d %H:%M:%S')
                
                created_at_str = created_at.strftime('%Y-%m-%d %H:%M:%S')
                index_date = '{}-{}-{}'.format(created_at_str[:4],created_at_str[5:7],created_at_str[8:10])
           
                
                      
                json_payload = {
                          "id" : pay_load['after']['id'],
                          "job_id" : pay_load['after']['job_id'],
                          "user_id" : pay_load['after']['user_id'],
                          "bpid" : pay_load['after']['bpid'],
                          "username" : pay_load['after']['username'],
                          "msisdn" : pay_load['after']['msisdn'],
                          "network" : pay_load['after']['network'],
                          "sender" : pay_load['after']['sender'],
                          "h_message" : pay_load['after']['h_message'],
                          "message" : pay_load['after']['message'],
                          "sms_count" : pay_load['after']['sms_count'],
                          "submit_date" : submit_date,
                          "delivery_date" : delivery_date,
                          "status" : pay_load['after']['status'],
                          "created_by" : pay_load['after']['created_by'],
                          "response" : pay_load['after']['response'],
                          "msgid" : pay_load['after']['msgid'],
                          "originated" : pay_load['after']['originated'],
                          "refid" : pay_load['after']['refid'],
                          "created_at" : created_at,
                          "updated_at" : updated_at,
                          "network_id" : pay_load['after']['network_id'],
                          "sms_price" : pay_load['after']['sms_price']
                         }
                
                new_list.append({
                      "_index": "deywuro-logs-{}".format(index_date), # The index on Elasticsearch
                      "_type": "index_type", # The document type
                      "_id": pay_load['after']['id'],
                      "_source":json_payload})
             

                if len(new_list) > 5000:
                    res=helpers.bulk(es, new_list)#,raise_on_error=False)
                    logging.info('success')
                    logging.info(f"[index_date:{index_date}] [submit_date:{submit_date}] [delivery_date:{delivery_date}]")
                    new_list = []


            except Exception as e:
                logging.info('#########################################################################')
                logging.info(f'failed with error {e}')
                logging.info('#########################################################################')
                continue
except Exception as e:
    logging.info('#########################################################################')
    logging.info(e)
    logging.info('#########################################################################')