package com.cmos.msgframe.admin.service.Impl;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Properties;
import java.util.Set;
import java.util.TreeMap;

import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;

import com.alibaba.rocketmq.client.ClientConfig;
import com.alibaba.rocketmq.client.impl.ClientRemotingProcessor;
import com.alibaba.rocketmq.client.impl.MQClientAPIImpl;
import com.alibaba.rocketmq.client.impl.factory.MQClientInstance;
import com.alibaba.rocketmq.common.MixAll;
import com.alibaba.rocketmq.common.protocol.RequestCode;
import com.alibaba.rocketmq.common.protocol.body.ClusterInfo;
import com.alibaba.rocketmq.common.protocol.body.KVTable;
import com.alibaba.rocketmq.common.protocol.route.BrokerData;
import com.alibaba.rocketmq.remoting.netty.NettyClientConfig;
import com.alibaba.rocketmq.remoting.protocol.RemotingCommand;
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.IBrokerService;
import com.cmos.msgframe.common.util.Constants;
import com.cmos.msgframe.common.util.OutPutParameter;

@Service("brokerService")
public class BrokerServiceImpl extends AbstractService implements
		IBrokerService {
	static final Logger logger = LoggerFactory
			.getLogger(TopicServiceImpl.class);

	@Override
	public OutPutParameter getBrokerList(String brokerAddr,String clusterNamep, int start, int end)
			throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		List<Map<String, String>> brokerList = new ArrayList<Map<String, String>>();
		OutPutParameter outData = new OutPutParameter();
		try {
			defaultMQAdminExt.start();
			ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt.examineBrokerClusterInfo();
			Iterator<Map.Entry<String, Set<String>>> itCluster = clusterInfoSerializeWrapper.getClusterAddrTable().entrySet().iterator();
		    while (itCluster.hasNext()){
		        Map.Entry<String, Set<String>> next = itCluster.next();
		        String clusterName = next.getKey();
		        Set<String> brokerNameSet = new HashSet<String>();
	            brokerNameSet.addAll(next.getValue());
	            for (String brokerName : brokerNameSet)
	            {
	            	Map<String, String> broMap = new HashMap<String, String>();
	            	broMap.put("clusterName", clusterName);
	            	broMap.put("brokerName", brokerName);
	            	 BrokerData brokerData = clusterInfoSerializeWrapper.getBrokerAddrTable().get(brokerName);
	            	 if (brokerData != null)
	                 {
	            		 Iterator<Map.Entry<Long, String>> itAddr = brokerData.getBrokerAddrs().entrySet().iterator();
	                     while (itAddr.hasNext())
	                     {
	                    	 Map.Entry<Long, String> next1 = itAddr.next();
	                    	  KVTable kvTable = defaultMQAdminExt.fetchBrokerRuntimeStats(next1.getValue());
	                    	  broMap.put("brokerVersion", kvTable.getTable().get("brokerVersionDesc"));
	                    	  broMap.put("brokerId", next1.getKey().toString());
	                    	  broMap.put("brokerAddr", next1.getValue());
	                     }
	                 }
	            	 brokerList.add(broMap);
	            }
	        }
		    if(StringUtils.isNotBlank(brokerAddr)&&!brokerAddr.equals("null")){
		    	brokerList = (List<Map<String,String>>)ListUtils.listFilerMap(brokerList, brokerAddr, false);
		    }
		    if(StringUtils.isNotBlank(clusterNamep)&&!clusterNamep.equals("null")){
		    	brokerList = (List<Map<String,String>>)ListUtils.listFilerMap(brokerList, clusterNamep, false);
		    }
		    Map<String, Object> totalMap = new HashMap<String, Object>();
		    totalMap.put("total", brokerList.size());
		    outData.setBean(totalMap);
		    outData.setBeans((List<Map<String, String>>) page(brokerList, start, end));
		    outData.setReturnCode(Constants.IS_OK);
		    outData.setReturnMessage("查询成功");
		    return outData;
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("查询出现异常");
		    logger.error("查询出现异常",e);
		    return outData;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}
	@Override
	public OutPutParameter getAllBroker()throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		List<Map<String, String>> brokerList = new ArrayList<Map<String, String>>();
		OutPutParameter outData = new OutPutParameter();
		try {
			defaultMQAdminExt.start();
			ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt.examineBrokerClusterInfo();
			Iterator<Map.Entry<String, Set<String>>> itCluster = clusterInfoSerializeWrapper.getClusterAddrTable().entrySet().iterator();
			while (itCluster.hasNext()){
				Map.Entry<String, Set<String>> next = itCluster.next();
				Set<String> brokerNameSet = new HashSet<String>();
				brokerNameSet.addAll(next.getValue());
				for (String brokerName : brokerNameSet)
				{
					Map<String, String> broMap = new HashMap<String, String>();
					broMap.put("name", brokerName);
					BrokerData brokerData = clusterInfoSerializeWrapper.getBrokerAddrTable().get(brokerName);
					if (brokerData != null)
					{
						Iterator<Map.Entry<Long, String>> itAddr = brokerData.getBrokerAddrs().entrySet().iterator();
						while (itAddr.hasNext())
						{
							Map.Entry<Long, String> next1 = itAddr.next();
							if(next1.getKey().toString().equals("0")){
								broMap.put("value", next1.getValue());
							}
						}
					}
					brokerList.add(broMap);
				}
			}
			outData.setBeans(brokerList);
			outData.setReturnCode(Constants.IS_OK);
			outData.setReturnMessage("查询成功");
			return outData;
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
			outData.setReturnMessage("查询出现异常");
			logger.error("查询出现异常",e);
			return outData;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}

	@Override
	public OutPutParameter getBrokerPropertyList(String brokerAddr,String property, int start, int end)
			throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		List broProList = new ArrayList();
		OutPutParameter outData = new OutPutParameter();
		ClientConfig clientConfig = new ClientConfig();
		MQClientInstance instance = new MQClientInstance(clientConfig,0,"");
		ClientRemotingProcessor processor = new ClientRemotingProcessor(instance);
		NettyClientConfig nettyClientConfig = new NettyClientConfig();
		MQClientAPIImpl mqClientAI = new MQClientAPIImpl(nettyClientConfig, processor);
		try {
			mqClientAI.start();
			defaultMQAdminExt.start();
			ClusterInfo clusterInfoSerializeWrapper = defaultMQAdminExt
					.examineBrokerClusterInfo();
			RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_BROKER_CONFIG, null);
			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) {
						if(StringUtils.isNotBlank(brokerAddr)&&!brokerAddr.equals("null")&&addr.indexOf(brokerAddr)==-1){
							continue;
						}
						RemotingCommand response = null;
					    response = mqClientAI.getRemotingClient().invokeSync(addr, request, 3000);
					    Properties properties = MixAll.string2Properties(new String(response.getBody(),MixAll.DEFAULT_CHARSET));
						if (null != properties) {
							Iterator<Entry<Object, Object>> it = properties.entrySet().iterator();
							while (it.hasNext()) {
								Map<Object, Object> proMap = new HashMap<Object, Object>();
								Entry<Object, Object> entry = it.next();
								proMap.put("key", entry.getKey());
								proMap.put("value", entry.getValue());
								proMap.put("brokerAddr", addr);
								proMap.put("clusterName", clusterN.toString());
								broProList.add(proMap);
							}
						}
					}
				}
			}
			if(StringUtils.isNotBlank(property)&&!property.equals("null")){
				broProList = (List<Map<String,String>>)ListUtils.listFilerMap(broProList, property, false);
			}
			Map<String, Object> totalMap = new HashMap<String, Object>();
			totalMap.put("total", broProList.size());
			outData.setBean(totalMap);
			outData.setBeans((List<Map<String, String>>) page(broProList, start, end));
			outData.setReturnCode(Constants.IS_OK);
		    outData.setReturnMessage("查询成功");
			return outData;
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("查询出现异常");
		    logger.error("查询出现异常",e);
		    return outData;
		} finally {
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
			mqClientAI.shutdown();
		}
	}

	@Override
	public OutPutParameter updateBrokerProperty(String brokerAddr, String clusterName,
			String key, String value) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			if(StringUtils.isEmpty(key)||key.equals("null")){
				outData.setReturnCode(Constants.SYSTEM_ERROR);
				outData.setReturnMessage("key值不能为空");
				return outData;
			}
			if(StringUtils.isEmpty(value)||value.equals("null")){
				outData.setReturnCode(Constants.SYSTEM_ERROR);
				outData.setReturnMessage("value值不能为空");
				return outData;
			}
			 Properties properties = new Properties();
             properties.put(key, value);
             if (StringUtils.isNotBlank(brokerAddr)&&!brokerAddr.equals("null")) {
                 defaultMQAdminExt.start();
                 defaultMQAdminExt.updateBrokerConfig(brokerAddr, properties);
                 outData.setReturnCode(Constants.IS_OK);
 				return outData;
             }else if (StringUtils.isNotBlank(clusterName)&&!clusterName.equals("null")) {
                 defaultMQAdminExt.start();
                 Set<String> masterSet =
                         CommandUtil.fetchMasterAddrByClusterName(defaultMQAdminExt, clusterName);
                 for (String tempBrokerAddr : masterSet) {
                     defaultMQAdminExt.updateBrokerConfig(tempBrokerAddr, properties);
                 }
                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.updateBrokerConfig(addr, properties);
        					}
        				}
        			}
    			outData.setReturnCode(Constants.IS_OK);
    			outData.setReturnMessage("修改成功");
 				return outData;
			}
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("查询出现异常");
		    logger.error("修改出现异常",e);
		    return outData;
		}finally{
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}

	@Override
	public OutPutParameter getBrokerDetail(String brokerAddr) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			if(StringUtils.isEmpty(brokerAddr)||brokerAddr.equals("null")){
				outData.setReturnCode(Constants.SYSTEM_ERROR);
				outData.setReturnMessage("brokerAddr值不能为空");
				return outData;
			}
			defaultMQAdminExt.start();
			KVTable kvTable = defaultMQAdminExt.fetchBrokerRuntimeStats(brokerAddr);
			Map<String, Object> tmp = new TreeMap<String, Object>();
			tmp.putAll(kvTable.getTable());
			outData.setBean(tmp);
			outData.setReturnCode(Constants.IS_OK);
			return outData;
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("获取broker资源详细信息出现异常");
		    logger.error("获取broker资源详细信息出现异常",e);
		    return outData;
		}finally{
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}
	@Override
	public OutPutParameter addBrokerProperty(String brokerAddr, String clusterName,
			String key, String value) throws Exception {
		DefaultMQAdminExt defaultMQAdminExt = getDefaultMQAdminExt();
		OutPutParameter outData = new OutPutParameter();
		try {
			if(StringUtils.isEmpty(key)||key.equals("null")){
				outData.setReturnCode(Constants.SYSTEM_ERROR);
				outData.setReturnMessage("key值不能为空");
				return outData;
			}
			if(StringUtils.isEmpty(value)||value.equals("null")){
				outData.setReturnCode(Constants.SYSTEM_ERROR);
				outData.setReturnMessage("value值不能为空");
				return outData;
			}
	
			 Properties properties = new Properties();
             properties.put(key, value);
             if (StringUtils.isNotBlank(brokerAddr)&&!brokerAddr.equals("null")) {
                 defaultMQAdminExt.start();
                 defaultMQAdminExt.updateBrokerConfig(brokerAddr, properties);
                 outData.setReturnCode(Constants.IS_OK);
  				return outData;
             }else if (StringUtils.isNotBlank(clusterName)&&!clusterName.equals("null")) {
                 defaultMQAdminExt.start();
                 Set<String> masterSet =
                         CommandUtil.fetchMasterAddrByClusterName(defaultMQAdminExt, clusterName);
                 for (String tempBrokerAddr : masterSet) {
                     defaultMQAdminExt.updateBrokerConfig(tempBrokerAddr, properties);
                 }
                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.updateBrokerConfig(addr, properties);
        					}
        				}
        			}
			outData.setReturnCode(Constants.IS_OK);
			return outData;
			}
		} catch (Exception e) {
			outData.setReturnCode(Constants.SYSTEM_ERROR);
		    outData.setReturnMessage("修改broker属性出现异常");
		    logger.error("修改broker属性出现异常",e);
		    return outData;
		}finally{
			shutdownDefaultMQAdminExt(defaultMQAdminExt);
		}
	}
}
