from confluent_kafka import Consumer
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_users.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': 'hellio-deywuro-users-0050',
  'auto.offset.reset': 'earliest'
})



try:
    consumer.subscribe(['hellio_live_c1.hellio.users'])
    new_list=[]
    now_time = datetime.now()
    es = Elasticsearch([str(elasticinfo['live'])])
    while True:
        try:
          if datetime.now() > now_time + timedelta(seconds=120):
            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 msg is None: 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']['full_name']:
                  full_name = "None"
                else:
                  full_name = pay_load['after']['full_name'].strip()

                if not pay_load['after']['username']:
                  username = "None"
                else:
                  username = pay_load['after']['username'].strip()
                
                if not pay_load['after']['organisation']:
                  organisation = "None"
                else:
                  organisation = pay_load['after']['organisation'].strip()
                
                if not pay_load['after']['status']:
                  status = "None"
                else:
                  status = pay_load['after']['status'].strip()
                
                if not pay_load['after']['location']:
                  location = "None"
                else:
                  location = pay_load['after']['location'].strip()

                if not pay_load['after']['phone_number']:
                  phone_number = "None"
                else:
                  phone_number = pay_load['after']['phone_number']

                if not pay_load['after']['email']:
                  email = "None"
                else:
                  email = pay_load['after']['email']

                if not pay_load['after']['created_at']:
                  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']['deleted_at']:
                  deleted_at = 'None'
                else:
                  deleted_at = parser.parse(pay_load['after']['deleted_at']).isoformat()
                  deleted_at = deleted_at.replace('T', ' ')
                  deleted_at = deleted_at.replace(deleted_at[19:], '')
                  updated_at = deleted_at.strptime(deleted_at, '%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])
                logging.info(f'index_date:{index_date}')
                      
                json_payload = {
                                  "id": pay_load['after']['id'],
                                  "full_name": str(full_name),
                                  "phone_number": str(phone_number), 
                                  "username": str(username),
                                  "email": str(email), 
                                  "role_id": pay_load['after']['role_id'], 
                                  "organisation": str(organisation),
                                  "created_at": created_at,
                                  "updated_at": updated_at,
                                  "deleted_at": deleted_at,
                                  "status": str(status),
                                  "location": str(location),
                                  "company_id" : pay_load['after']['company_id'],
                                  "services": "None",
                                  "code":"None",

                                }
                
                new_list.append({
                      "_index": "hellio-deywuro-users-{}".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('created')
                    new_list = []


            except Exception as e:
                logging.info('#########################################################################')
                logging.info(e)
                logging.info(' ')
                logging.info('failed')
                logging.info('#########################################################################')
                continue
except Exception as e:
    logging.info('#########################################################################')
    logging.info(e)
    logging.info('#########################################################################')