package com.cmos.msgframe.admin.service.Impl;

import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Set;

import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;

import com.alibaba.rocketmq.common.MixAll;
import com.alibaba.rocketmq.common.UtilAll;
import com.alibaba.rocketmq.common.admin.ConsumeStats;
import com.alibaba.rocketmq.common.admin.OffsetWrapper;
import com.alibaba.rocketmq.common.message.MessageQueue;
import com.alibaba.rocketmq.common.protocol.body.ClusterInfo;
import com.alibaba.rocketmq.common.protocol.body.ConsumerConnection;
import com.alibaba.rocketmq.common.protocol.body.GroupList;
import com.alibaba.rocketmq.common.protocol.body.TopicList;
import com.alibaba.rocketmq.common.subscription.SubscriptionGroupConfig;
import com.alibaba.rocketmq.tools.admin.DefaultMQAdminExt;
import com.alibaba.rocketmq.tools.command.CommandUtil;
import com.cmos.msgframe.admin.common.AbstractService;
import com.cmos.msgframe.admin.service.IConsumerService;
import com.cmos.msgframe.common.util.Constants;
import com.cmos.msgframe.common.util.OutPutParameter;

@Service("consumerService")
public class ConsumerServiceImpl extends AbstractService implements
		IConsumerService {
	static final Logger logger = LoggerFactory
			.getLogger(ConsumerServiceImpl.class);

	@Override
	public OutPutParameter getConsumerList(String consumerGroup, String topic,
			int start, int end) throws Exception {
		
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		List<Map<String,String>> resultList = new ArrayList<Map<String,String>>();
		OutPutParameter outData = new OutPutParameter();
		try {
			if (StringUtils.isNotBlank(consumerGroup)&&!consumerGroup.equals("null")) {
				defaultMQAdminExt.start();
				ConsumeStats consumeStats = null;
				try {
					consumeStats = defaultMQAdminExt.examineConsumeStats(consumerGroup);
				} catch (Exception e) {
					outData.setBeans(resultList);
					outData.setReturnCode(Constants.SYSTEM_ERROR);
					outData.setReturnMessage("获取消费组连接列表出现异常");
					return outData;
				}
				resultList.add(getConsumerByGroup(defaultMQAdminExt,consumeStats, consumerGroup));
				Map<String, Object> totalMap = new HashMap<String, Object>();
				totalMap.put("total", resultList.size());
				outData.setBean(totalMap);
				outData.setBeans(resultList);
				return outData;
			}else if(StringUtils.isNotBlank(topic)&&!topic.equals("null")){
				defaultMQAdminExt.start();
				GroupList groupList = defaultMQAdminExt.queryTopicConsumeByWho(topic);
				for (String consumerGroupTemp:groupList.getGroupList()) {
					ConsumeStats consumeStats = null;
					try {
						consumeStats = defaultMQAdminExt.examineConsumeStats(consumerGroupTemp);
					} catch (Exception e) {
						logger.error("消费组"+consumerGroupTemp+"没有进行消息消费");
						continue;
					}
					resultList.add(getConsumerByGroup(defaultMQAdminExt,consumeStats, consumerGroupTemp));
				}
				Map<String, Object> totalMap = new HashMap<String, Object>();
				totalMap.put("total", resultList.size());
				outData.setBean(totalMap);
				outData.setBeans(resultList);
				return outData;
			}else {
				defaultMQAdminExt.start();
				TopicList topicList = defaultMQAdminExt.fetchAllTopicList();
				for (String topicTemp : topicList.getTopicList()) {
					if (topicTemp.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
						continue;
					}
					GroupList groupList = defaultMQAdminExt.queryTopicConsumeByWho(topicTemp);
					for (String consumerGroupTemp:groupList.getGroupList()) {
						ConsumeStats consumeStats = null;
						try {
							consumeStats = defaultMQAdminExt.examineConsumeStats(consumerGroupTemp);
						} catch (Exception e) {
							logger.error("消费组"+consumerGroupTemp+"没有进行消息消费");
							continue;
						}
						resultList.add(getConsumerByGroup(defaultMQAdminExt,consumeStats, consumerGroupTemp));
					}
				}
				Map<String, Object> totalMap = new HashMap<String, Object>();
				totalMap.put("total", resultList.size());
				outData.setBean(totalMap);
				outData.setBeans(resultList);
				outData.setReturnCode(Constants.IS_OK);
				return outData;
			}
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("查询消费组信息失败");
		    logger.error("查询消费组信息失败",e);
		    return outData;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}

	private Map<String, String> getConsumerByGroup(
			DefaultMQAdminExt defaultMQAdminExt, ConsumeStats consumeStats,
			String consumerGroup) {
		Map<String, String> map = new HashMap<String, String>();
		map.put("consumerGroup", consumerGroup);
		List<MessageQueue> mqList = new LinkedList<MessageQueue>();
		mqList.addAll(consumeStats.getOffsetTable().keySet());
		Collections.sort(mqList);
		String topicTemp = "";
		String clusterNameTemp = "";
		Map clusterSet = null;
		try {
			ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt
					.examineBrokerClusterInfo();
			clusterSet = clusterInfoSerializeWrapper.getClusterAddrTable();
		} catch (Exception e) {
			logger.error("获取主题所属集群出现异常");
		}
		for (MessageQueue mq : mqList) {
			if(mq.getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)){
				continue;
			}
			topicTemp = frontStringAtLeast(mq.getTopic(), 32);
			if (clusterSet.size() > 0) {
				for (Object clusterN : clusterSet.keySet()) {
					clusterNameTemp = clusterN.toString();
				}
			}
			break;
		}
		map.put("topic", topicTemp);
		map.put("clusterName", clusterNameTemp);
		map.put("diffTotal", String.valueOf(consumeStats.computeTotalDiff()));
		map.put("consumerTps", String.valueOf(consumeStats.getConsumeTps()));
		ConsumerConnection cc = null;
		try {
			cc = defaultMQAdminExt.examineConsumerConnectionInfo(consumerGroup);
		} catch (Exception e) {
			logger.error("消费者不在线");
		}
		if (cc != null) {
			map.put("consumerCount",String.valueOf(cc.getConnectionSet().size()));
			map.put("MessageModel", cc.getMessageModel().toString());
		} else {
			map.put("consumerCount","0");
			map.put("MessageModel", "UNKOWN");
		}
		return map;
	}

	@Override
	public OutPutParameter addSubgroup(String brokerAddr, String clusterName,
			String groupName, String consumeEnable,
			String consumeFromMinEnable, String consumeBroadcastEnable,
			String retryQueueNums, String retryMaxTimes, String brokerId)
			throws Exception {
	       DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
	       OutPutParameter outData = new OutPutParameter();
	       try {
	            SubscriptionGroupConfig subscriptionGroupConfig = new SubscriptionGroupConfig();
	            subscriptionGroupConfig.setConsumeBroadcastEnable(false);
	            subscriptionGroupConfig.setConsumeFromMinEnable(false);
	            if(StringUtils.isEmpty(groupName)||groupName.equals("null")){
	            	outData.setReturnCode(Constants.SYSTEM_ERROR);
	            	outData.setReturnMessage("订阅主题不能为空");
	            	return outData;
	            }
	            // groupName
	            subscriptionGroupConfig.setGroupName(groupName);

	            // consumeEnable
	            if (StringUtils.isNotBlank(consumeEnable)&&consumeEnable.equals("null")) {
	                subscriptionGroupConfig.setConsumeEnable(Boolean.parseBoolean(consumeEnable.trim()));
	            }

	            // consumeFromMinEnable
	            if (StringUtils.isNotBlank(consumeFromMinEnable)) {
	                subscriptionGroupConfig.setConsumeFromMinEnable(Boolean.parseBoolean(consumeFromMinEnable
	                    .trim()));
	            }

	            // consumeBroadcastEnable
	            if (StringUtils.isNotBlank(consumeBroadcastEnable)) {
	                subscriptionGroupConfig.setConsumeBroadcastEnable(Boolean.parseBoolean(consumeBroadcastEnable
	                    .trim()));
	            }

	            // retryQueueNums
	            if (StringUtils.isNotBlank(retryQueueNums)) {
	                subscriptionGroupConfig.setRetryQueueNums(Integer.parseInt(retryQueueNums.trim()));
	            }

	            // retryMaxTimes
	            if (StringUtils.isNotBlank(retryMaxTimes)) {
	                subscriptionGroupConfig.setRetryMaxTimes(Integer.parseInt(retryMaxTimes.trim()));
	            }

	            // brokerId
	            if (StringUtils.isNotBlank(brokerId)) {
	                subscriptionGroupConfig.setBrokerId(Long.parseLong(brokerId.trim()));
	            }

	            if (StringUtils.isNotBlank(brokerAddr)) {
	                defaultMQAdminExt.start();
	                defaultMQAdminExt.createAndUpdateSubscriptionGroupConfig(brokerAddr, subscriptionGroupConfig);
	                outData.setReturnCode(Constants.IS_OK);
					return outData;

	            } else if (StringUtils.isNotBlank(clusterName)) {
	                defaultMQAdminExt.start();
	                Set<String> masterSet =CommandUtil.fetchMasterAddrByClusterName(defaultMQAdminExt, clusterName);
	                for (String addr : masterSet) {
	                    defaultMQAdminExt.createAndUpdateSubscriptionGroupConfig(addr, subscriptionGroupConfig);
	                }
	                outData.setReturnCode(Constants.IS_OK);
					return outData;
	            }
	            else {
	            	defaultMQAdminExt.start();
	            	ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt.examineBrokerClusterInfo();
        			Map clusterSet = clusterInfoSerializeWrapper.getClusterAddrTable();
        			if (clusterSet.size() > 0) {
        				for (Object clusterN : clusterSet.keySet()) {
        					 Set<String> masterSet =CommandUtil.fetchMasterAddrByClusterName(defaultMQAdminExt, clusterN.toString());
    		                 for (String addr : masterSet) {
    		                    defaultMQAdminExt.createAndUpdateSubscriptionGroupConfig(addr, subscriptionGroupConfig);
    		                 }
        				}
        			}
        			outData.setReturnCode(Constants.IS_OK);
    				return outData;
	            }
	       }catch (Exception e) {
				outData.setReturnCode(Constants.SYSTEM_ERROR);
			    outData.setReturnMessage("添加消费组信息出现异常");
			    logger.error("查添加消费组信息出现异常",e);
			    return outData;
			} finally {
				shutdownDefaultMQAdminExt(defaultMQAdminExt);
			}
	}

	@Override
	public OutPutParameter updateSubgroup(String brokerAddr, String clusterName,
			String groupName, String consumeEnable,
			String consumeFromMinEnable, String consumeBroadcastEnable,
			String retryQueueNums, String retryMaxTimes, String brokerId)
			throws Exception {
		return  addSubgroup(brokerAddr,clusterName,groupName,consumeEnable,consumeFromMinEnable,consumeBroadcastEnable,retryQueueNums,retryMaxTimes,brokerId);
	}

	@Override
	public OutPutParameter deleteSubgroup(String clusterName, String groupName)
			throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			if(StringUtils.isNotBlank(clusterName)){
				defaultMQAdminExt.start();
			    Set<String> masterSet = CommandUtil.fetchMasterAddrByClusterName(defaultMQAdminExt, clusterName);
                for (String master : masterSet) {
                	defaultMQAdminExt.deleteSubscriptionGroup(master, groupName);
                }
                outData.setReturnCode(Constants.IS_OK);
			}
			return outData;
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("删除消费组信息出现异常");
		    logger.error("删除消费组信息出现异常",e);
		    return outData;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}

	@Override
	public OutPutParameter getConsumerStats(String consumerGroup, int start,
			int end) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		List<Map<String,String>> resultList = new ArrayList<Map<String,String>>();
		OutPutParameter outData = new OutPutParameter();
		try {
			defaultMQAdminExt.start();
			if (StringUtils.isNotBlank(consumerGroup)) {
				ConsumeStats consumeStats = defaultMQAdminExt.examineConsumeStats(consumerGroup);
				List<MessageQueue> mqList = new LinkedList<MessageQueue>();
				mqList.addAll(consumeStats.getOffsetTable().keySet());
				Collections.sort(mqList);
				for (MessageQueue mq : mqList) {
					Map<String, String> map = new HashMap<String, String>();
					 OffsetWrapper offsetWrapper = consumeStats.getOffsetTable().get(mq);
					 long diff = offsetWrapper.getBrokerOffset() - offsetWrapper.getConsumerOffset();
					String topic = UtilAll.frontStringAtLeast(mq.getTopic(), 32);
					String brokerName = UtilAll.frontStringAtLeast(mq.getBrokerName(), 32);
					String queueId = String.valueOf(mq.getQueueId());
					String brokerOffset = String.valueOf(offsetWrapper.getBrokerOffset());
					String consumerOffset = String.valueOf(offsetWrapper.getConsumerOffset());
					String needConsumer = String.valueOf(diff);
					map.put("topic", topic);
					map.put("brokerName", brokerName);
					map.put("queueId", queueId);
					map.put("brokerOffset", brokerOffset);
					map.put("consumerOffset", consumerOffset);
					map.put("needConsumer", needConsumer);
					resultList.add(map);
				}
				Map<String, Object> totalMap = new HashMap<String, Object>();
				totalMap.put("total", resultList.size());
				outData.setBean(totalMap);
				outData.setBeans((List<Map<String, String>>) page(resultList, start, end));
				outData.setReturnCode(Constants.IS_OK);
			}
			return outData;
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("查询消费者详细信息出现异常");
		    logger.error("查询消费者详细信息出现异常",e);
		    return outData;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}
}
