消息系统

This commit is contained in:
jiangs 2026-09-01 14:31:36 +08:00
parent e3c6eae37e
commit 2e77bcb10c
14 changed files with 244 additions and 36 deletions

View File

@ -2,7 +2,9 @@ package com.czblood.bus.mapper;
import com.czblood.bus.pojo.XkMsgDetail;
import java.util.List;
public interface XkMsgDetailMapper {
XkMsgDetail queryOne(int msgId);
List<XkMsgDetail> queryOne(int msgId);
void insert(XkMsgDetail msgDetail);
}

View File

@ -2,7 +2,12 @@ package com.czblood.bus.mapper;
import com.czblood.bus.pojo.XkMsgMain;
import java.util.List;
public interface XkMsgMainMapper {
XkMsgMain queryOne(int msgId);
void insert(XkMsgMain msgMain);
List<XkMsgMain> queryList(XkMsgMain msgMain);
void update(XkMsgMain msgMain);
int queryNoReadCount();
}

View File

@ -24,7 +24,7 @@ import java.util.List;
import java.util.stream.Collectors;
/**
* 登录接口
* 外部共用接口
*/
@Component
public class BloodExtInterfaceService implements ExtSystemInterface {

View File

@ -6,27 +6,29 @@ import lombok.Getter;
public enum MsgType {
//取血通知消息
BLOOD_OUT_NOTICE("BLOOD_OUT_NOTICE"),
BLOOD_OUT_NOTICE("BLOOD_OUT_NOTICE","取血通知消息"),
//输血申请单保存消息
APPLY_NOTICE("APPLY_NOTICE"),
APPLY_NOTICE("APPLY_NOTICE","输血申请单保存消息"),
/** 单发消息 */
SINGLE("SINGLE"),
SINGLE("SINGLE","单发消息"),
/** 群发消息 */
// BROADCAST("BROADCAST"),
/** 心跳 */
HEARTBEAT("HEARTBEAT"),
HEARTBEAT("HEARTBEAT","心跳"),
/** 系统通知 */
NOTICE("NOTICE");
NOTICE("NOTICE","系统通知");
private final String code;
private final String desc;
MsgType(String code) {
MsgType(String code,String desc) {
this.code = code;
this.desc = desc;
}
public static MsgType fromCode(String code) {

View File

@ -46,8 +46,8 @@ public class SendMessageUtil {
xkMsgMainTemp.setRead_Flag("N"); //是否已读
xkMsgMainTemp.setSend_Flag("Y"); //是否已发送
xkMsgMainTemp.setAccept_Flag("N"); //是否已接收
xkMsgMainTemp.setMsg_Title("输血申请单保存");
xkMsgMainTemp.setMsg_Type(MsgType.APPLY_NOTICE.getCode()); //发血消息
xkMsgMainTemp.setMsg_Title(MsgType.APPLY_NOTICE.getDesc()); //申请单消息
xkMsgMainTemp.setMsg_Type(MsgType.APPLY_NOTICE.getCode());
xkMsgMainTemp.setWrite_Date(currentTimeDate);
//写消息数据明细表
@ -61,6 +61,7 @@ public class SendMessageUtil {
@Transactional(rollbackFor = Exception.class)
public void saveData(XkMsgMain xkMsgMain, XkMsgDetail xkMsgDetail) {
xkMsgMainMapper.insert(xkMsgMain);
xkMsgDetail.setMsgid(xkMsgMain.getMsgid());
xkMsgDetailMapper.insert(xkMsgDetail);
}

View File

@ -2,12 +2,11 @@ package com.czblood.bus.websocket;
import cn.hutool.json.JSONObject;
import cn.hutool.json.JSONUtil;
import com.czblood.bus.pojo.XkMsgDetail;
import com.czblood.bus.pojo.XkMsgMain;
import com.czblood.bus.pojo.websocket.WebSocketParam;
import com.czblood.bus.utils.BloodBankUtil;
import com.czblood.common.core.domain.entity.SysDept;
import com.czblood.common.core.domain.entity.SysUser;
import com.czblood.common.utils.HosIdHolder;
import com.czblood.common.utils.spring.SpringUtils;
import com.czblood.system.mapper.SysDeptMapper;
import lombok.extern.slf4j.Slf4j;
@ -59,11 +58,24 @@ public class WebSocketServer {
/**
* 收到客户端消息后调用的方法
* 根据 JSON 中的 type 字段判断:单发 / 群发 / 心跳 / 通知
* // ★ 手动设置 hos_id 到 ThreadLocal
* // WebSocket 线程没有 HTTP 请求上下文和 SecurityContext,HosIdHolder 无法自动获取 hos_id
* // 因此从 onOpen 时缓存的 WebSocketParam.loginUser 中取出 hos_id 手动注入
*/
@OnMessage
public void onMessage(String message, @PathParam("userId") String userId) {
log.info("收到客户端[{}]消息: {}", userId, message);
String hosId = null;
WebSocketParam senderParam = onlineUsers.get(userId);
if (senderParam != null && senderParam.getLoginUser() != null) {
hosId = senderParam.getLoginUser().getHos_id();
}
if (hosId != null && !hosId.isEmpty()) {
HosIdHolder.setHosId(hosId);
log.info("WebSocket 手动注入 hos_id={}, userId={}", hosId, userId);
} else {
log.warn("WebSocket 未能获取到发送者的 hos_id, userId={}", userId);
}
try {
// 解析 JSON 为 WsMessage
WsMessage msg = JSONUtil.toBean(message, WsMessage.class);
@ -76,7 +88,7 @@ public class WebSocketServer {
switch (msg.getType()) {
case SINGLE:
// 单发:发给指定用户
sendMsg(msg.getTo(), message);
sendMsg(msg.getTo(), msg);
break;
case BLOOD_OUT_NOTICE:
case APPLY_NOTICE:
@ -97,6 +109,9 @@ public class WebSocketServer {
} catch (Exception e) {
e.printStackTrace();
log.error("消息处理异常 userId={}", userId, e);
} finally {
// ★ 清理 ThreadLocal,防止 Tomcat 线程池复用导致脏数据残留
HosIdHolder.clear();
}
}
@ -142,16 +157,18 @@ public class WebSocketServer {
/**
* 单发:向指定 userId 所有在线 ws 会话推送消息
*/
public static void sendMsg(String userId, String content) {
public static void sendMsg(String userId, WsMessage msg) {
WebSocketParam param = onlineUsers.get(userId);
if (param == null || param.getSessionSet().isEmpty()) {
log.info("没有用户在线,无法发送消息。。。。userId={}",userId);
return;
}
String jsonStr = JSONUtil.toJsonStr(msg);
for (Session wsSession : param.getSessionSet()) {
if (wsSession.isOpen()) {
try {
wsSession.getBasicRemote().sendText(content);
log.info("发消息成功 userId={}: {}", userId, content);
wsSession.getBasicRemote().sendText(jsonStr);
log.info("发消息成功 userId={}: {}", userId, jsonStr);
} catch (Exception e) {
log.error("单发消息失败 userId={}: {}", userId, e.getMessage());
}
@ -166,8 +183,7 @@ public class WebSocketServer {
public static void sendToSomeOne(WsMessage msg){
MsgType type = msg.getType();
String content = msg.getContent();
String userId = msg.getTo();
WebSocketParam param = onlineUsers.get(userId);
String billNo = "";
try{
JSONObject data = (JSONObject)msg.getData();
@ -177,18 +193,9 @@ public class WebSocketServer {
return;
}
if (param == null || param.getSessionSet().isEmpty()) {
return;
}
for (Session wsSession : param.getSessionSet()) {
if (wsSession.isOpen()) {
try {
wsSession.getBasicRemote().sendText(content);
} catch (Exception e) {
log.error("单发消息失败 userId={}: {}", userId, e.getMessage());
}
}
}
//保存消息数据,写一条消息记录
SpringUtils.getBean(SendMessageUtil.class).saveMessage(billNo,content);
switch (type) {
//输血申请单保存通知
case APPLY_NOTICE:
@ -202,13 +209,12 @@ public class WebSocketServer {
sysUserQuery.setDeptId(deptId);
List<SysUser> sysUserList = SpringUtils.getBean(BloodBankUtil.class).queryUserList(sysUserQuery);
for (SysUser sysUser : sysUserList) {
sendMsg(sysUser.getUserName(), content);
sendMsg(sysUser.getUserName(), msg);
}
}
break;
}
//保存消息数据,写一条消息记录
SpringUtils.getBean(SendMessageUtil.class).saveMessage(billNo,content);
}
}

View File

@ -28,8 +28,34 @@
<if test="msg_Title != null"> #{msg_Title},</if>
#{write_Date})
</insert>
<update id="update">
update xk_msg_main set
<if test="accept_Person != null"> accept_Person = #{accept_Person},</if>
<if test="accept_Date != null"> accept_Date = #{accept_Date},</if>
<if test="read_Flag != null"> read_Flag = #{read_Flag},</if>
<if test="send_Flag != null"> send_Flag = #{send_Flag},</if>
<if test="accept_Flag != null"> accept_Flag = #{accept_Flag},</if>
where msgid = #{msgid}
</update>
<select id="queryOne" resultType="com.czblood.bus.pojo.XkMsgMain">
select * from xk_msg_main where msgid = #{msgId}
</select>
<select id="queryList" resultType="com.czblood.bus.pojo.XkMsgMain">
select * from xk_msg_main
<where>
<if test="accept_Dept != null">and accept_Dept = #{accept_Dept}</if>
<if test="send_Person != null">and send_Person = #{send_Person}</if>
<if test="accept_Person != null">and accept_Person = #{accept_Person}</if>
<if test="read_Flag != null">and read_Flag = #{read_Flag}</if>
<if test="send_Flag != null">and send_Flag = #{send_Flag}</if>
<if test="accept_Flag != null">and accept_Flag = #{accept_Flag}</if>
<if test="msg_Type != null">and msg_Type = #{msg_Type}</if>
<if test="msg_Kind != null">and msg_Kind = #{msg_Kind}</if>
<if test="msgid != null">and msgid = #{msgid}</if>
</where>
</select>
<select id="queryNoReadCount" resultType="java.lang.Integer">
select count(*) as noReadCount from xk_msg_main where isnull(read_Flag, 'N') = 'N'
</select>
</mapper>

View File

@ -0,0 +1,52 @@
package com.czblood.busBloodbank.controller;
import com.czblood.bus.pojo.XkMsgMain;
import com.czblood.busBloodbank.pojo.VO.MsgVO;
import com.czblood.busBloodbank.service.MsgService;
import com.czblood.common.core.controller.BaseController;
import com.czblood.common.core.domain.Result;
import com.czblood.common.core.page.TableDataInfo;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiModelProperty;
import io.swagger.annotations.ApiOperation;
import org.springframework.web.bind.annotation.*;
import javax.annotation.Resource;
import java.util.List;
@Api(tags = "消息系统")
@RestController
@RequestMapping("/bus/bloodbank/msg")
public class MsgController extends BaseController {
@Resource
MsgService msgService;
@ApiOperation("查询消息列表")
@PostMapping("/queryMsgList")
public TableDataInfo queryMsgList(@RequestBody XkMsgMain xkMsgMain){
startPage();
List<XkMsgMain> msgMainList = msgService.queryMsgList(xkMsgMain);
return getDataTable(msgMainList);
}
@ApiOperation("获取某个消息详情")
@GetMapping("/queryOne")
public Result queryOne(int msgId){
return msgService.queryOne(msgId);
}
@ApiModelProperty("获取未读消息数量")
@GetMapping("/queryNoReadCount")
public Result queryNoReadCount(){
return msgService.queryNoReadCount();
}
@ApiOperation("更新消息状态")
@GetMapping("/updateFlag")
public Result updateFlag(MsgVO msgVO){
return msgService.updateFlag(msgVO);
}
}

View File

@ -0,0 +1,17 @@
package com.czblood.busBloodbank.pojo.DTO;
import com.czblood.bus.pojo.XkMsgDetail;
import com.czblood.bus.pojo.XkMsgMain;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.List;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class MsgDTO {
private XkMsgMain xkMsgMain;
private List<XkMsgDetail> xkMsgDetailList;
}

View File

@ -0,0 +1,20 @@
package com.czblood.busBloodbank.pojo.VO;
import io.swagger.annotations.ApiModelProperty;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.Date;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class MsgVO {
@ApiModelProperty("消息ID")
private Integer msgId;
@ApiModelProperty("接收人")
private String accept_Person;
@ApiModelProperty("接收时间")
private Date accept_Date;
}

View File

@ -0,0 +1,14 @@
package com.czblood.busBloodbank.service;
import com.czblood.bus.pojo.XkMsgMain;
import com.czblood.busBloodbank.pojo.VO.MsgVO;
import com.czblood.common.core.domain.Result;
import java.util.List;
public interface MsgService {
List<XkMsgMain> queryMsgList(XkMsgMain xkMsgMain);
Result queryOne(int msgId);
Result updateFlag(MsgVO msgVO);
Result queryNoReadCount();
}

View File

@ -11,6 +11,7 @@ import com.czblood.bus.utils.BloodInterfaceUtil;
import com.czblood.bus.utils.ReportUtil;
import com.czblood.bus.websocket.MsgType;
import com.czblood.bus.websocket.WebSocketServer;
import com.czblood.bus.websocket.WsMessage;
import com.czblood.busBloodbank.mapper.BloodOutMapper;
import com.czblood.busBloodbank.pojo.DTO.BloodOutDTO;
import com.czblood.busBloodbank.pojo.DTO.BloodOutInfoDTO;
@ -482,7 +483,7 @@ public class BloodOutServiceImpl implements BloodOutService {
xkMsgMainTemp.setRead_Flag("N"); //是否已读
xkMsgMainTemp.setSend_Flag("Y"); //是否已发送
xkMsgMainTemp.setAccept_Flag("N"); //是否已接收
xkMsgMainTemp.setMsg_Title("取血通知");
xkMsgMainTemp.setMsg_Title(MsgType.BLOOD_OUT_NOTICE.getDesc());
xkMsgMainTemp.setMsg_Type(MsgType.BLOOD_OUT_NOTICE.getCode()); //发血消息
xkMsgMainTemp.setWrite_Date(currentTimeDate);
@ -499,7 +500,8 @@ public class BloodOutServiceImpl implements BloodOutService {
List<SysUser> sysUserList = bloodBankUtil.queryUserList(sysUserQuery);
for (SysUser sysUser : sysUserList) {
String userName = sysUser.getUserName();
WebSocketServer.sendMsg(userName,msgBody);
WsMessage msg = new WsMessage();
WebSocketServer.sendMsg(userName,msg);
}
return new Result("0","操作成功");
}

View File

@ -0,0 +1,60 @@
package com.czblood.busBloodbank.service.impl;
import com.czblood.bus.mapper.XkMsgDetailMapper;
import com.czblood.bus.mapper.XkMsgMainMapper;
import com.czblood.bus.pojo.XkMsgDetail;
import com.czblood.bus.pojo.XkMsgMain;
import com.czblood.bus.utils.BloodBankUtil;
import com.czblood.busBloodbank.pojo.DTO.MsgDTO;
import com.czblood.busBloodbank.pojo.VO.MsgVO;
import com.czblood.busBloodbank.service.MsgService;
import com.czblood.common.core.domain.Result;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.util.List;
@Service
public class MsgServiceImpl implements MsgService {
@Resource
XkMsgMainMapper xkMsgMainMapper;
@Resource
XkMsgDetailMapper xkMsgDetailMapper;
@Override
public List<XkMsgMain> queryMsgList(XkMsgMain xkMsgMain) {
return xkMsgMainMapper.queryList(xkMsgMain);
}
@Override
public Result queryOne(int msgId) {
XkMsgMain xkMsgMain = xkMsgMainMapper.queryOne(msgId);
List<XkMsgDetail> xkMsgDetailList = xkMsgDetailMapper.queryOne(msgId);
MsgDTO msgDTO = new MsgDTO();
msgDTO.setXkMsgMain(xkMsgMain);
msgDTO.setXkMsgDetailList(xkMsgDetailList);
return new Result("0","查询成功",msgDTO);
}
@Override
public Result updateFlag(MsgVO msgVO) {
XkMsgMain xkMsgMainOne = xkMsgMainMapper.queryOne(msgVO.getMsgId());
if(xkMsgMainOne == null) return new Result("-1","没有找到该消息记录!");
//只记录第一次读取人和时间
if(BloodBankUtil.isBank(xkMsgMainOne.getRead_Flag()).equals("Y")) return new Result("0","该消息已读取!");
XkMsgMain xkMsgMain = new XkMsgMain();
xkMsgMain.setMsgid(msgVO.getMsgId());
xkMsgMain.setAccept_Person(msgVO.getAccept_Person());
xkMsgMain.setAccept_Date(msgVO.getAccept_Date());
xkMsgMain.setRead_Flag("Y");
xkMsgMainMapper.update(xkMsgMain);
return new Result("0","更新成功");
}
@Override
public Result queryNoReadCount() {
int noReadCount = xkMsgMainMapper.queryNoReadCount();
return new Result("0","查询成功",noReadCount);
}
}

View File

@ -33,6 +33,7 @@ import java.util.Properties;
args = {MappedStatement.class, Object.class, RowBounds.class, ResultHandler.class,
CacheKey.class, BoundSql.class})
})
//没有用到===========================================
public class HosIdInterceptor implements Interceptor {
private static final String PARAM_NAME = "hos_id";