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
from elasticsearch import Elasticsearch, helpers
import re




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)


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-users-0000',
  'auto.offset.reset': 'earliest'
})


try:
    data = []
    id = 0
    now_time = datetime.now()
    es = Elasticsearch(['http://127.0.0.1:9200/'])
    consumer.subscribe(['ai.kannel.credits'])
    while True:
        try:
          if datetime.now() > now_time + timedelta(seconds=120):
            now_time = datetime.now()
            es.indices.refresh(index="hellio-users-{}".format(datetime.today().strftime('%Y')))
            logging.info('index refreshed')
          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'))
              id += 1

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

              if pay_load['after']['balance'] == None:
                balance = "None"
              else:
                balance = pay_load['after']['balance']

              created_at = datetime.now().isoformat()
              updated_at = datetime.now().isoformat() 
              

              
              json_payload = {
                            "id": id,
                            "full_name": "None",
                            "phone_number": "None", 
                            "username": str(username),
                            "email": "None", 
                            "role_id": "None", 
                            "organisation": "None",
                            "created_at": created_at,
                            "updated_at": updated_at,
                            "status": "None",
                            "location": "None",
                            "company_id" : "None",
                            "balance": str(balance)
              }

              res = es.index(index="hellio-users-{}".format(datetime.today().strftime('%Y')), id="kannel-{}".format(id), body=json_payload)
              logging.info(res)
          except Exception as e:
              logging.info('#########################################################################')
              logging.info(e)
              logging.info(' ')
              logging.info('error')
              logging.info('#########################################################################')
              continue
except Exception as e:
    logging.info('#########################################################################')
    logging.info(e)
    logging.info('#########################################################################')

