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/credits_histories_consumer.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': 'user_credit_histories-0001',
  'auto.offset.reset': 'earliest'
})



try:
    consumer.subscribe(['hellio_live.hellio.credit_histories'])
    new_list=[]
    now_time = datetime.now()
    es = Elasticsearch([str(elasticinfo['live'])])
    while True:
        try:
          if datetime.now() > now_time + timedelta(seconds=60):
            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']:
                  created_at = "None"
                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 = "None"
                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:], '')
                  deleted_at = datetime.strptime(deleted_at, '%Y-%m-%d %H:%M:%S')
                
                if not pay_load['after']['credited_by']:
                    credited_by = 'None'
                else:
                    credited_by = pay_load['after']['credited_by']

                if not pay_load['after']['date_credited']:
                  date_credited = 'None'     
                else:
                  date_credited = datetime.utcfromtimestamp((pay_load['after']['date_credited'])/1000).strftime('%Y-%m-%d %H:%M:%S')
                  date_credited = datetime.strptime(date_credited, '%Y-%m-%d %H:%M:%S')
                
                
                json_payload = {
                          "id" : pay_load['after']['id'],
                          "user_id" : pay_load['after']['user_id'],
                          "username" : pay_load['after']['username'],
                          "price" : pay_load['after']['price'],
                          "credit_alloted" : pay_load['after']['credit_alloted'],
                          "old_bal" : pay_load['after']['old_bal'],
                          "new_bal" : pay_load['after']['new_bal'],
                          "credit_used" : pay_load['after']['credit_used'],
                          "credited" : pay_load['after']['credited'],
                          "currency" : pay_load['after']['currency'],
                          "credited_by" : credited_by,
                          "comments" : pay_load['after']['comments'],
                          "status" : pay_load['after']['status'],
                          "date_credited" : date_credited,
                          "created_at" : created_at,
                          "updated_at" : updated_at,
                          "deleted_at" : deleted_at
                         }
                
                new_list.append({
                      "_index": "hellio-credit-histories", # 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')
                    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('#########################################################################')