zhongrj
2023-07-10 943e5e1fc85c696593ea3b81f7f95e235477b3c0
站内信信息写入,sse连接,断开修改
7 files modified
3 files added
273 ■■■■ changed files
skjcmanager/skjcmanager-service-api/skjcmanager-alerts-api/src/main/java/cn/gistack/alerts/notice/feign/NoticeClient.java 1 ●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/entity/MessageRecord.java 2 ●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/feign/IMessageClient.java 43 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/feign/IMessageClientFallback.java 25 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/vo/MessageRecordVO.java 13 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/AsyncNoticeHandle.java 32 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/NoticeHandle.java 93 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/NoticeStrategyImpl.java 22 ●●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/controller/SSEController.java 41 ●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/server/SSEServer.java 1 ●●●● patch | view | raw | blame | history
skjcmanager/skjcmanager-service-api/skjcmanager-alerts-api/src/main/java/cn/gistack/alerts/notice/feign/NoticeClient.java
@@ -18,7 +18,6 @@
    /**
     * 发送告警通知
     * @param options 单独超时时间设置
     * @param name
     * @param content
     * @param type
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/entity/MessageRecord.java
@@ -37,7 +37,7 @@
    private String content;
    /**
     * 状态  已读,未读
     * 状态 0未读  1已读
     */
    private Integer status;
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/feign/IMessageClient.java
New file
@@ -0,0 +1,43 @@
package cn.gistack.sm.message.feign;
import cn.gistack.sm.message.entity.MessageRecord;
import org.springframework.cloud.openfeign.FeignClient;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestParam;
import java.util.List;
/**
 * Feign失败配置
 * @author zhongrj
 * @date 2023-07-07
 */
@FeignClient(
    value = "blade-sm",
    fallback = IMessageClientFallback.class
)
public interface IMessageClient {
    String API_PREFIX = "/client";
    String SAVE_MESSAGE_RECORD = API_PREFIX + "/saveMessageRecord";
    String GET_NOT_READ_NUMBER = API_PREFIX + "/getNotReadNumber";
    /**
     * 保存站内信记录信息
     * @param messageRecord
     * @return
     */
    @PostMapping(SAVE_MESSAGE_RECORD)
    void saveMessageRecord(@RequestBody List<MessageRecord> messageRecord);
    /**
     * 查询用户未读消息数量
     * @param phone
     * @return
     */
    @GetMapping(GET_NOT_READ_NUMBER)
    Integer getNotReadNumber(@RequestParam("phone") String phone);
}
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/feign/IMessageClientFallback.java
New file
@@ -0,0 +1,25 @@
package cn.gistack.sm.message.feign;
import cn.gistack.sm.message.entity.MessageRecord;
import org.springframework.stereotype.Component;
import java.util.List;
/**
 * Feign失败配置
 * @author zhongrj
 * @date 2023-07-07
 */
@Component
public class IMessageClientFallback implements IMessageClient {
    @Override
    public void saveMessageRecord(List<MessageRecord> messageRecord) {
    }
    @Override
    public Integer getNotReadNumber(String phone) {
        return null;
    }
}
skjcmanager/skjcmanager-service-api/skjcmanager-sm-api/src/main/java/cn/gistack/sm/message/vo/MessageRecordVO.java
New file
@@ -0,0 +1,13 @@
package cn.gistack.sm.message.vo;
import cn.gistack.sm.message.entity.MessageRecord;
import lombok.Data;
@Data
public class MessageRecordVO extends MessageRecord {
    /**
     * 是否分页(导出时不分页)
     */
    private Integer isPage;
}
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/AsyncNoticeHandle.java
@@ -27,6 +27,38 @@
    /**
     * 水位接近汛限告警
     * @param templateId
     * @param uuid
     * @param alarmRule
     * @param adCode
     */
    @Async
    public void waterComeOverHandle(String templateId, String uuid, AlarmRule alarmRule, String adCode) {
        // 调用中台服务接口查询数据
        JSONArray ztData = noticeHandle.getZtData("", ZtApiUrlConstant.res_diff_dead_z_api);
        // 只发送站内信
        SmsRequestTemplate smsRequestTemplate = noticeHandle.getSendSmsTemplate(ztData, templateId,alarmRule, ZtApiDataColumnConstant.waterLessDeadList,true,null);
        // 拼接内容后面的联系方式,当前模板县水利部人员姓名电话
        // 遍历,拼接县水利部人员电话
        for (List<String> list : smsRequestTemplate.getTemplateContent()) {
            // 后面依次拼接 巡查责任人姓名,技术责任人姓名,市县水利部门人员姓名电话(暂时不要)
            noticeHandle.addPatrolTechnologyInfo(list);
            // 市县水利部门人员姓名电话(暂时不要)
//                list.add("测试:13112341234");
            // 保存告警记录信息
//            AlarmRecord alarmRecord = noticeHandle.saveAlarmRecord(list,alarmRule,uuid);
//            if (null != alarmRecord) {
//                list.add(alarmRecord.getId().toString());
                    list.add("999");
//            }
        }
        // 保存站内信记录信息
        noticeHandle.saveInStationInfo(smsRequestTemplate,alarmRule);
    }
    /**
     * 水位初次超汛限告警 处理
     * @param templateId
     * @param uuid
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/NoticeHandle.java
@@ -1,6 +1,8 @@
package cn.gistack.alerts.notice.service.impl;
import cn.com.flaginfo.sdk.cmc.api.sms.send.SMSSendRequest;
import cn.gistack.alerts.alarmRule.entity.AlarmRecord;
import cn.gistack.alerts.alarmRule.entity.AlarmRecordDetail;
import cn.gistack.alerts.alarmRule.entity.AlarmRule;
import cn.gistack.alerts.alarmRule.service.AlarmRecordService;
import cn.gistack.alerts.notice.constant.ZtApiUrlConstant;
@@ -10,14 +12,19 @@
import cn.gistack.alerts.sms.entity.SmsTemplate;
import cn.gistack.alerts.sms.service.ISmsTemplateService;
import cn.gistack.alerts.sms.vo.SmsRequestTemplate;
import cn.gistack.alerts.sse.server.SSEServer;
import cn.gistack.common.utils.SpringContextUtil;
import cn.gistack.sm.message.entity.MessageRecord;
import cn.gistack.sm.message.feign.IMessageClient;
import cn.gistack.sm.sjztmd.feign.IAttResManagePersonClient;
import cn.gistack.sm.sjztmd.vo.PersonVO;
import cn.gistack.system.feign.ISysClient;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.alibaba.nacos.shaded.com.google.gson.JsonObject;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
import org.springblade.core.secure.utils.AuthUtil;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpEntity;
import org.springframework.http.HttpHeaders;
@@ -47,6 +54,9 @@
    @Autowired
    private RestTemplate restTemplate;
    @Autowired
    private IMessageClient messageClient;
@@ -420,4 +430,87 @@
        }
        return null;
    }
    /**
     * 保存站内信记录信息
     * @param smsRequestTemplate
     * @param alarmRule
     * @return
     */
    public Map<String, Object> saveInStationInfo(SmsRequestTemplate smsRequestTemplate, AlarmRule alarmRule) {
        // 取出数据
        SMSSendRequest request = new SMSSendRequest();
        SmsTemplate smsTemplate = smsRequestTemplate.getSmsTemplate();
        request.setTemplateId(smsTemplate.getTemplateId());
        List<List<String>> templateContent = smsRequestTemplate.getTemplateContent();
        // 先保存结果记录
        List<AlarmRecordDetail> alarmRecordDetailList = new ArrayList<>();
        List<MessageRecord> messageRecordList = new ArrayList<>();
        Set<String> phoneSet = new HashSet<>();
        // 遍历处理
        for (List<String> list : templateContent) {
            // 设置站内信记录信息
            MessageRecord messageRecord = new MessageRecord();
            messageRecord.setTheme(alarmRule.getRuleName());
            messageRecord.setSender(AuthUtil.getUserId().toString());
            messageRecord.setSource("智能告警");
            messageRecord.setStatus(2);
            // 设置告警记录信息
            AlarmRecordDetail alarmRecordDetail = new AlarmRecordDetail();
            alarmRecordDetail.setReservoirNumber(list.get(0));
            alarmRecordDetail.setAlarmRecordId(Long.parseLong(list.get(list.size() - 1)));
            alarmRecordDetail.setAlarmMode("站内信");
            // 未读
            alarmRecordDetail.setStatus(1);
            // 删除 最后一个 告警记录id
            list.remove(list.size() - 1);
            // 删除 0 水库编码
            list.remove(0);
            alarmRecordDetail.setPhone(list.get(0));
            // 设置手机号
            request.setUserNumber(list.get(0));
            messageRecord.setRecipient(list.get(0));
            phoneSet.add(list.get(0));
            // 内容处理
            String s = smsTemplate.getContent().replaceAll("\\{.+?\\}", "%s");
            // 需要删除最前面的手机号
            list.remove(0);
            // 格式转换
            String format = String.format(s, list.toArray());
            // 设置内容
            request.setMessageContent(format);
            messageRecord.setContent(format);
            // 保存记录
            alarmRecordDetail.setAlarmContent(format);
            alarmRecordDetail.setCreateTime(new Date());
            // 放入集合
            alarmRecordDetailList.add(alarmRecordDetail);
            messageRecordList.add(messageRecord);
        }
        // 保存告警记录详情信息
        saveMessageRecord(messageRecordList);
        Map<String, Object> map = new HashMap<>(2);
        // 返回
        return map;
    }
    /**
     * 保存站内信记录信息
     * @param messageRecordList
     */
    public void saveMessageRecord(List<MessageRecord> messageRecordList) {
        // 保存记录
        messageClient.saveMessageRecord(messageRecordList);
        // 发起通知
        // 遍历
        for (MessageRecord messageRecord : messageRecordList) {
            Map<String, Object> map = new HashMap<>(2);
            // 查询当前用户未读消息个数
            map.put("count",messageClient.getNotReadNumber(messageRecord.getRecipient()));
            map.put("record",messageRecord);
            SSEServer.sendMessage("web:"+messageRecord.getRecipient(),new JSONObject(map).toJSONString());
        }
    }
}
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/notice/service/impl/NoticeStrategyImpl.java
@@ -196,31 +196,13 @@
     * @return
     */
    public String waterComeOverHandle(String arg1, String arg2) {
        String adCode = null;
        String templateId = "znxsw10000001";
        String uuid = templateId + UUID.randomUUID().toString();
        // 查询当前策略对应的告警规则信息
        AlarmRule alarmRule = alarmRuleService.getOne(new QueryWrapper<AlarmRule>().eq("rule_name", arg1));
        if (null!=alarmRule) {
            // 调用中台服务接口查询数据
            JSONArray ztData = noticeHandle.getZtData("", ZtApiUrlConstant.res_diff_dead_z_api);
            // 只发送站内信
            SmsRequestTemplate smsRequestTemplate = noticeHandle.getSendSmsTemplate(ztData, templateId,alarmRule, ZtApiDataColumnConstant.waterLessDeadList,true,null);
            // 拼接内容后面的联系方式,当前模板县水利部人员姓名电话
            // 遍历,拼接县水利部人员电话
            for (List<String> list : smsRequestTemplate.getTemplateContent()) {
                // 后面依次拼接 巡查责任人姓名,技术责任人姓名,市县水利部门人员姓名电话(暂时不要)
                noticeHandle.addPatrolTechnologyInfo(list);
                // 市县水利部门人员姓名电话(暂时不要)
//                list.add("测试:13112341234");
                // 保存告警记录信息
                AlarmRecord alarmRecord = noticeHandle.saveAlarmRecord(list,alarmRule,uuid);
                if (null != alarmRecord) {
                    list.add(alarmRecord.getId().toString());
//                    list.add("999");
                }
            }
            // 保存站内信记录信息
//            smsService.saveInStationInfo(smsRequestTemplate);
            asyncNoticeHandle.waterComeOverHandle(templateId,uuid,alarmRule, null);
        }
        return "水位接近汛限告警";
    }
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/controller/SSEController.java
@@ -1,16 +1,25 @@
package cn.gistack.alerts.sse.controller;
import cn.gistack.alerts.notice.service.impl.NoticeHandle;
import cn.gistack.alerts.sse.server.SSEServer;
import cn.gistack.alerts.sse.vo.SseVO;
import cn.gistack.sm.message.entity.MessageRecord;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.util.ArrayList;
import java.util.List;
@Slf4j
@RestController
@CrossOrigin
@RequestMapping("/sse/sse")
@AllArgsConstructor
public class SSEController {
    private final NoticeHandle noticeHandle;
    /**
     * 建立连接
@@ -19,7 +28,19 @@
     */
    @GetMapping("/connect")
    public SseEmitter connect(SseVO sse){
        return SSEServer.connect(sse.getType() + ":" + sse.getUserId());
        String userId = sse.getType() + ":" + sse.getUserId();
        return SSEServer.connect(userId);
    }
    /**
     * 断开连接
     * @param sse
     * @return
     */
    @GetMapping("/disconnect")
    public void disconnect(SseVO sse){
        String userId = sse.getType() + ":" + sse.getUserId();
        SSEServer.removeUser(userId);
    }
    /**
@@ -27,15 +48,13 @@
     * @throws InterruptedException
     */
    @GetMapping("/process")
    public void sendMessage() throws InterruptedException {
        SSEServer.sendMessage("web:123","hello!");
//        for(int i=0; i<=100; i++){
//            if(i>50&&i<70){
//                Thread.sleep(500L);
//            }else{
//                Thread.sleep(100L);
//            }
//            SSEServer.batchSendMessage(String.valueOf(i));
//        }
    public void sendMessage(){
        MessageRecord messageRecord = new MessageRecord();
        messageRecord.setRecipient("1642816241683865602");
        messageRecord.setTheme("测试123");
        messageRecord.setStatus(0);
        List<MessageRecord> list = new ArrayList<>();
        list.add(messageRecord);
        noticeHandle.saveMessageRecord(list);
    }
}
skjcmanager/skjcmanager-service/skjcmanager-alerts/src/main/java/cn/gistack/alerts/sse/server/SSEServer.java
@@ -34,6 +34,7 @@
        //数量+1
        count.getAndIncrement();
        log.info("create new sse connect ,current user:{}",userId);
        log.info("count",count.getAndIncrement());
        return sseEmitter;
    }
    /**