package slite.lib_java;

import java.time.Duration;
import java.util.Arrays;
import java.util.HashMap;

import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.clients.producer.*;

public abstract class SliteKafkaClient<K, V>
{
	KafkaConsumer<K,V> consumer;
	KafkaProducer<K,V> producer;
	HashMap<String, String> settings;
	public String clientId;
	String group = null;
	String servers;
	String offset = "earliest";
	String[] subscribeTopics;
	String defaultProduceTopic = "default";
	int subSecTimeout = 5;
	protected String keySerializer = "ByteArray";
	protected String valueSerializer = "ByteArray";

	public abstract void receive(String topic, K key, long timestampMillis, V data) throws Exception;

	/**
	 * Will get all of the config settings from environment variables
	 */
	public SliteKafkaClient() throws Exception // Will get all of the config settings from environment variables
	{
	}

	/**
	 * Will use these settings as defaults, will override them with environment variables
	 */
	public SliteKafkaClient(String defaultId, String defaultGroup, String defaultServers, String defaultOffset, String defaultProduceTopic, String[] defaultSubscribeTopics) throws Exception
	{
		this.clientId = defaultId;
		this.group = defaultGroup;
		this.servers = defaultServers;
		this.offset = defaultOffset;
		this.subscribeTopics = defaultSubscribeTopics;
		this.defaultProduceTopic = defaultProduceTopic;

		this.init();
	}

	protected void init() throws Exception
	{
		this.loadEnv();
		this.initSettings();
		this.initConsumer();
		this.initProducer();
	}

	private void loadEnv() throws Exception
	{
		String tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_CLIENT_ID");
		if(tmp!=null && !tmp.isEmpty()) this.clientId = tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_SERVERS");
		if(tmp!=null && !tmp.isEmpty()) this.servers = tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_OFFSET");
		if(tmp!=null && !tmp.isEmpty()) this.offset = tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_SUB_GROUP");
		if(tmp!=null && !tmp.isEmpty()) this.group = tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_SUB_TOPICS");
		if(tmp!=null && !tmp.isEmpty()) this.subscribeTopics = tmp.split(";");
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_PUB_TOPIC");
		if(tmp!=null && !tmp.isEmpty()) this.defaultProduceTopic = tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_SERIALIZER_KEY");
		if(tmp!=null && !tmp.isEmpty()) this.keySerializer = tmp;
		
		tmp = System.getenv("SLITE_LIB_JAVA_KAFKA_SERIALIZER_VALUE");
		if(tmp!=null && !tmp.isEmpty()) this.valueSerializer = tmp;
	}

	private void initSettings() throws Exception
	{
		final SliteKafkaClient<K, V> thisObj = this;
		
		this.settings = new HashMap<String, String>()
		{
			{
				put("client.id", thisObj.clientId+"");
				if(thisObj.group!=null) put("group.id", thisObj.group);
				put("bootstrap.servers", thisObj.servers);
				if(thisObj.offset!=null) put("auto.offset.reset", thisObj.offset);
				put("key.deserializer", "org.apache.kafka.common.serialization."+thisObj.keySerializer+"Deserializer");
				put("value.deserializer", "org.apache.kafka.common.serialization."+thisObj.valueSerializer+"Deserializer");
				put("key.serializer", "org.apache.kafka.common.serialization."+thisObj.keySerializer+"Serializer");
				put("value.serializer", "org.apache.kafka.common.serialization."+thisObj.valueSerializer+"Serializer");

				System.out.println(this.toString().replace(",", "\n"));
			}
		};
	}

	private void initProducer() throws Exception
	{
		this.producer = new KafkaProducer(this.settings);
	}

	private void initConsumer() throws Exception
	{
		this.consumer = new KafkaConsumer(this.settings);
		if(this.subscribeTopics!=null) 
		{
			this.consumer.subscribe(Arrays.asList(this.subscribeTopics));
			final SliteKafkaClient<K, V> thisObj = this;

			new Thread()
			{
				public void run()
				{
					while(true)
					{
						ConsumerRecords<K, V> records = thisObj.consumer.poll(Duration.ofSeconds(thisObj.subSecTimeout));
						if(!records.isEmpty())
						{
							records.forEach
							(
								(ConsumerRecord<K, V> record) -> 
								{
									var topic = record.topic();
									K key = record.key();
									var timestampMillis = record.timestamp();
									V data = record.value();
									try
									{
										thisObj.receive(topic, key, timestampMillis, data);
									}
									catch (Exception e)
									{
										e.printStackTrace();
									}
								}
							);
							thisObj.consumer.commitSync(Duration.ofSeconds(5));
						}
					}
				}
			}.start();
		}
	}

	/**
	 * Only use this method if your key is of type String. Any other key types will throw a cast exception.
	 */
	public void produce(V data) throws Exception
	{
		this.produce(data, (K)this.clientId, this.defaultProduceTopic);
	}

	public void produce(V data, K key) throws Exception
	{
		this.produce(data, key, this.defaultProduceTopic);
	}

	public void produce(V data, K key, String topic) throws Exception
	{
		this.produce(data, key, topic, System.currentTimeMillis());
	}

	public void produce(V data, K key, String topic, Long timestamp) throws Exception
	{
		this.producer.send(new ProducerRecord<K,V>(topic, null, timestamp, key, data));
	}
}