package com.cmos.msgframe.admin.service.Impl;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Set;

import oracle.net.aso.b;

import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;

import com.alibaba.rocketmq.common.TopicConfig;
import com.alibaba.rocketmq.common.admin.TopicOffset;
import com.alibaba.rocketmq.common.admin.TopicStatsTable;
import com.alibaba.rocketmq.common.message.MessageQueue;
import com.alibaba.rocketmq.common.protocol.body.ClusterInfo;
import com.alibaba.rocketmq.common.protocol.body.TopicList;
import com.alibaba.rocketmq.common.protocol.route.BrokerData;
import com.alibaba.rocketmq.common.protocol.route.QueueData;
import com.alibaba.rocketmq.common.protocol.route.TopicRouteData;
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.common.ListUtils;
import com.cmos.msgframe.admin.service.ITopicService;
import com.cmos.msgframe.common.util.Constants;
import com.cmos.msgframe.common.util.OutPutParameter;

@Service("topicService")
public class TopicServiceImpl extends AbstractService implements ITopicService {
	static final Logger logger = LoggerFactory
			.getLogger(TopicServiceImpl.class); 

	public OutPutParameter getTopicList(final String topicCode, int start,
			int end) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		List<Map<String, String>> resultList = new ArrayList<Map<String, String>>();
		OutPutParameter outPutParameter = new OutPutParameter();
		Map clusterSet = null;
		try {
			defaultMQAdminExt.start();
			ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt
					.examineBrokerClusterInfo();
			clusterSet = clusterInfoSerializeWrapper.getClusterAddrTable();
			
			TopicList topicList = defaultMQAdminExt.fetchAllTopicList();
			Map<String, Object> totalMap = new HashMap<String, Object>();
			
			for (String topicName : topicList.getTopicList()) {
				Map<String, String> map = new HashMap<String, String>();
				map.put("topicName", topicName);
				if (clusterSet.size() > 0) {
					for (Object clusterN : clusterSet.keySet()) {
						map.put("ClusterName", clusterN.toString());
						break;
					}
				}
				List<QueueData> queList = defaultMQAdminExt
						.examineTopicRouteInfo(topicName).getQueueDatas();
				for (QueueData queueData : queList) {
					map.put("readNums", queueData.getReadQueueNums() + "");
					map.put("writeNums", queueData.getWriteQueueNums() + "");
					break;
				}
				TopicStatsTable topicStatsTable = defaultMQAdminExt.examineTopicStats(topicName);
				List<MessageQueue> mqList = new LinkedList<MessageQueue>();
				mqList.addAll(topicStatsTable.getOffsetTable().keySet());
				Collections.sort(mqList);
				Long maxOffset = 0L;
				for (MessageQueue mq : mqList) {
					TopicOffset topicOffset = topicStatsTable.getOffsetTable().get(mq);
					maxOffset = (maxOffset+topicOffset.getMaxOffset());
				}
				map.put("MaxOffset", maxOffset+"");
				resultList.add(map);
			}

			if (StringUtils.isNotBlank(topicCode)&&!topicCode.equals("null") && resultList.size() > 0) {
				resultList = (List<Map<String, String>>) ListUtils
						.listFilerMap(resultList, topicCode, false);
			}
			totalMap.put("total", resultList.size());
			outPutParameter.setBean(totalMap);
			outPutParameter.setBeans((List<Map<String, String>>) page(
					resultList, start, end));
			outPutParameter.setReturnCode(Constants.IS_OK);
			return outPutParameter;
		} catch (Exception e) {
			outPutParameter.setReturnCode(Constants.SYSTEM_ERROR);
			outPutParameter.setReturnMessage("没有获取到主题信息或查询出现异常");
			logger.error("没有获取到主题信息或查询出现异常",e);
			return outPutParameter;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}

	}

	@Override
	public OutPutParameter addTopic(String topic, String readQueueNums,
			String writeQueueNums, String clusterName)
			throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			TopicConfig topicConfig = new TopicConfig();
			if(StringUtils.isEmpty(topic)||topic.equals("null")){
            	outData.setReturnCode(Constants.SYSTEM_ERROR);
            	outData.setReturnMessage("主题不能为空");
            	return outData;
            }
			topicConfig.setTopicName(topic);
			topicConfig.setReadQueueNums(8);
			topicConfig.setWriteQueueNums(8);
			topicConfig.setOrder(false);
			if (StringUtils.isNotBlank(readQueueNums)&&!readQueueNums.equals("null")) {
				topicConfig.setReadQueueNums(Integer.parseInt(readQueueNums));
			}

			if (StringUtils.isNotBlank(writeQueueNums)&&!writeQueueNums.equals("null")) {
				topicConfig.setWriteQueueNums(Integer.parseInt(writeQueueNums));
			}

			 if (StringUtils.isNotBlank(clusterName)&&!clusterName.equals("null")) {

				defaultMQAdminExt.start();

				Set<String> masterSet = CommandUtil
						.fetchMasterAddrByClusterName(defaultMQAdminExt,
								clusterName);
				for (String addr : masterSet) {
					defaultMQAdminExt.createAndUpdateTopicConfig(addr,
							topicConfig);
				}
				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.createAndUpdateTopicConfig(addr,
									topicConfig);
						}
					}
				}
				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 updateTopic(String topic, String readQueueNums,
			String writeQueueNums, String clusterName) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			if(StringUtils.isEmpty(topic)||topic.equals("null")){
				outData.setReturnCode(Constants.SYSTEM_ERROR);
				outData.setReturnMessage("主题不能为空");
				return outData;
			}
			TopicConfig topicConfig = new TopicConfig();
			topicConfig.setTopicName(topic);
			topicConfig.setReadQueueNums(8);
			topicConfig.setWriteQueueNums(8);
			topicConfig.setOrder(false);
			if (StringUtils.isNotBlank(readQueueNums)&&!readQueueNums.equals("null")) {
				topicConfig.setReadQueueNums(Integer.parseInt(readQueueNums));
			}

			if (StringUtils.isNotBlank(writeQueueNums)&&!writeQueueNums.equals("null")) {
				topicConfig.setWriteQueueNums(Integer.parseInt(writeQueueNums));
			}

			if (StringUtils.isNotBlank(clusterName)&&!clusterName.equals("null")) {

				defaultMQAdminExt.start();

				Set<String> masterSet = CommandUtil
						.fetchMasterAddrByClusterName(defaultMQAdminExt,
								clusterName);
				for (String addr : masterSet) {
					defaultMQAdminExt.createAndUpdateTopicConfig(addr,
							topicConfig);
				}
				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.createAndUpdateTopicConfig(addr,
									topicConfig);
						}
					}
				}
				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 deleteTopic(String topic, String clusterName)
			throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			if(StringUtils.isEmpty(topic)||topic.equals("null")){
            	outData.setReturnCode(Constants.SYSTEM_ERROR);
            	outData.setReturnMessage("主题不能为空");
            	return outData;
            }
			if (StringUtils.isNotBlank(clusterName)&&!clusterName.equals("null")) {
				defaultMQAdminExt.start();
				Set<String> masterSet = CommandUtil
						.fetchMasterAddrByClusterName(defaultMQAdminExt,
								clusterName);
				defaultMQAdminExt.deleteTopicInBroker(masterSet, topic);
				Set<String> nameServerSet = null;
				if (StringUtils.isNotBlank(configureInitializer
						.getNamesrvAddr())) {
					String[] ns = configureInitializer.getNamesrvAddr().split(
							";");
					nameServerSet = new HashSet<String>(Arrays.asList(ns));
				}
				defaultMQAdminExt.deleteTopicInNameServer(nameServerSet, topic);
				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());
						defaultMQAdminExt.deleteTopicInBroker(masterSet, topic);
						Set<String> nameServerSet = null;
						if (StringUtils.isNotBlank(configureInitializer
								.getNamesrvAddr())) {
							String[] ns = configureInitializer.getNamesrvAddr()
									.split(";");
							nameServerSet = new HashSet<String>(
									Arrays.asList(ns));
						}
						defaultMQAdminExt.deleteTopicInNameServer(
								nameServerSet, topic);
						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);
		}
		outData.setReturnCode(Constants.SYSTEM_ERROR);
		return outData;
	}

	@Override
	public OutPutParameter getTopicDetail(String topicCode) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		Map<String, Object> detailMap = new HashMap<String, Object>();
		try {
			defaultMQAdminExt.start();
			ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt
					.examineBrokerClusterInfo();
			Map clusterSet = clusterInfoSerializeWrapper.getClusterAddrTable();
			if(StringUtils.isEmpty(topicCode)||topicCode.equals("null")){
				throw new Exception("主题名称不能为空");
			}
			detailMap.put("topicCode", topicCode);
//			Set<String> clusterSet = defaultMQAdminExt
//					.getTopicClusterList(topicCode);
//			Iterator<String> it = clusterSet.iterator();
//			while (it.hasNext()) {
//				String clusterName = (String) it.next();
//				detailMap.put("ClusterName", clusterName);
//				break;
//			}
			if (clusterSet.size() > 0) {
				for (Object clusterN : clusterSet.keySet()) {
					detailMap.put("ClusterName", clusterN.toString());
					break;
				}
			}
			List<QueueData> queList = defaultMQAdminExt
					.examineTopicRouteInfo(topicCode).getQueueDatas();
			for (QueueData queueData : queList) {
				detailMap.put("readNums", queueData.getReadQueueNums() + "");
				detailMap.put("writeNums", queueData.getWriteQueueNums() + "");
				break;
			}
			outData.setBean(detailMap);
			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 getTopicDetailById(String topicCode, int start, int end)
			throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		List<Map<String, String>> resultList = new ArrayList<Map<String, String>>();
		try {
			if (StringUtils.isEmpty(topicCode)||topicCode.equals("null")) {
				throw new Exception("主题名称不能为空");
			}

			defaultMQAdminExt.start();
			TopicStatsTable topicStatsTable = defaultMQAdminExt.examineTopicStats(topicCode);
			TopicRouteData topicRouteData = defaultMQAdminExt.examineTopicRouteInfo(topicCode);
			List<MessageQueue> mqList = new LinkedList<MessageQueue>();
			mqList.addAll(topicStatsTable.getOffsetTable().keySet());
			Collections.sort(mqList);
			for (MessageQueue mq : mqList) {
				Map<String, String> detailMap = new HashMap<String, String>();
				detailMap.put("topicCode", topicCode);
				TopicOffset topicOffset = topicStatsTable.getOffsetTable().get(mq);
				String humanTimestamp = "";
				if (topicOffset.getLastUpdateTimestamp() > 0) {
					humanTimestamp = timeMillisToHumanString2(topicOffset.getLastUpdateTimestamp());
				}
				detailMap.put("LastUpdateTime", humanTimestamp);
				List<BrokerData> RouteDataList = topicRouteData.getBrokerDatas();
				for (BrokerData data : RouteDataList) {
					Map<Long, String> brokerMap = data.getBrokerAddrs();
					if (null != brokerMap.get(0L)&& data.getBrokerName().equals(mq.getBrokerName())) {
						detailMap.put("brokerName", mq.getBrokerName());
						detailMap.put("brokerId", "0");
						detailMap.put("brokerAddr", brokerMap.get(0));
					}
					detailMap.put("MinOffset", topicOffset.getMinOffset() + "");
					detailMap.put("MaxOffset", topicOffset.getMaxOffset() + "");
					detailMap.put("QueueId", mq.getQueueId() + "");
				}
				resultList.add(detailMap);
			}
			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);
		}
	}
}
