zhuyifan

feat(server): 实现 WebSocket 服务并添加消息推送功能

......@@ -43,6 +43,7 @@ dependencies {
implementation "org.apache.poi:poi-ooxml:4.1.2" //poi
implementation 'cn.hutool:hutool-all:5.8.27' //huTool-util
implementation 'com.fasterxml.jackson.dataformat:jackson-dataformat-xml:2.11.2' // Jackson XML 模块
implementation 'org.springframework.boot:spring-boot-starter-websocket'// websocket
implementation 'com.aliyun:aliyun-java-sdk-core:4.6.4' //阿里云SDK核心库
implementation 'com.aliyun:aliyun-java-sdk-dysmsapi:2.2.1'//阿里云短信服务SDK
......
......@@ -4,7 +4,7 @@ import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.web.servlet.error.ErrorMvcAutoConfiguration;
@SpringBootApplication(exclude = ErrorMvcAutoConfiguration.class)
@SpringBootApplication(exclude = ErrorMvcAutoConfiguration.class, scanBasePackages = {"com.infoloop.tianting"})
public class App {
public static void main(String[] args) {
......
package com.infoloop.tianting.config;
import static springfox.documentation.builders.RequestHandlerSelectors.basePackage;
import com.github.xiaoymin.knife4j.spring.annotations.EnableKnife4j;
import com.github.xiaoymin.knife4j.spring.extension.OpenApiExtensionResolver;
import com.infoloop.tianting.constant.CommonConstants;
import io.swagger.models.auth.In;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import springfox.documentation.builders.ApiInfoBuilder;
......@@ -43,7 +39,7 @@ public class SwaggerConfig {
"/v2/api-docs",
"/v3/api-docs",
"/v3/api-docs/**",
"/doc.html",
"/doc.html"
};
/*引入Knife4j提供的扩展类*/
......
package com.infoloop.tianting.config;
import com.infoloop.tianting.intercepter.WebSocketHandshakeInterceptor;
import com.infoloop.tianting.server.WebSocketServer;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.config.annotation.EnableWebSocket;
import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry;
import org.springframework.web.socket.server.standard.ServletServerContainerFactoryBean;
@Slf4j
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(webSocketServer(), "/ws/message")
.setAllowedOrigins("*")
.addInterceptors(new WebSocketHandshakeInterceptor())
;
}
@Bean
public ServletServerContainerFactoryBean createWebSocketContainer() {
ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean();
container.setMaxTextMessageBufferSize(8192);
container.setMaxBinaryMessageBufferSize(8192);
return container;
}
@Bean
public WebSocketServer webSocketServer() {
return new WebSocketServer();
}
}
......@@ -58,4 +58,8 @@ public interface CommonConstants {
String OPEN_ID = "openId";
String USER_TYPE = "userType";
String USER_ID = "userId";
}
......
package com.infoloop.tianting.enums;
import com.infoloop.tianting.server.session.UserTypeEnum;
import lombok.AllArgsConstructor;
import lombok.Getter;
@Getter
@AllArgsConstructor
public enum LoginSourceEnum{
CUSTOMER(1, "住院客户"),
MINI_PROGRAM(2, "阳光厅客户"),
OPERATOR(3, "院区端"),
KDS(4, "KDS"),
CUSTOMER(1, UserTypeEnum.MICRO),
MINI_PROGRAM(2, UserTypeEnum.MICRO),
OPERATOR(3, UserTypeEnum.CLIENT),
KDS(4, UserTypeEnum.KDS),
;
/**
......@@ -19,5 +20,14 @@ public enum LoginSourceEnum{
/**
* 订单类型的名字
*/
private final String name;
private final UserTypeEnum userType;
public static LoginSourceEnum getByValue(Integer value) {
for (LoginSourceEnum orderTypeEnum : LoginSourceEnum.values()) {
if (orderTypeEnum.getValue().equals(value)) {
return orderTypeEnum;
}
}
return null;
}
}
......
package com.infoloop.tianting.intercepter;
import cn.dev33.satoken.stp.StpUtil;
import cn.hutool.core.convert.Convert;
import com.infoloop.tianting.constant.CommonConstants;
import com.infoloop.tianting.enums.LoginSourceEnum;
import org.springframework.http.server.ServerHttpRequest;
import org.springframework.http.server.ServerHttpResponse;
import org.springframework.util.StringUtils;
import org.springframework.web.socket.WebSocketHandler;
import org.springframework.web.socket.server.HandshakeInterceptor;
import java.util.Map;
@SuppressWarnings("all")
public class WebSocketHandshakeInterceptor implements HandshakeInterceptor {
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) {
final var authorization = request.getHeaders().getFirst(CommonConstants.AUTHORIZATION);
if (StringUtils.isEmpty(authorization)) {
return false;
}
final var loginSource = Convert.toInt(StpUtil.getExtra(CommonConstants.LOGIN_SOURCE));
final var loginSourceEnum = LoginSourceEnum.getByValue(loginSource);
attributes.put(CommonConstants.USER_TYPE, loginSourceEnum.getUserType().name());
attributes.put(CommonConstants.USER_ID, StpUtil.getLoginIdAsString());
return true;
}
@Override
public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception exception) {
}
}
package com.infoloop.tianting.server;
import com.infoloop.tianting.constant.CommonConstants;
import com.infoloop.tianting.server.message.SocketMessage;
import com.infoloop.tianting.server.session.UserSessionData;
import com.infoloop.tianting.server.session.UserSessionKey;
import com.infoloop.tianting.server.session.UserTypeEnum;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.socket.CloseStatus;
import org.springframework.web.socket.TextMessage;
import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.handler.TextWebSocketHandler;
import java.io.IOException;
import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@Slf4j
public class WebSocketServer extends TextWebSocketHandler {
private static final ConcurrentHashMap<UserSessionKey, UserSessionData> userSessions = new ConcurrentHashMap<>();
private static final long INACTIVITY_TIMEOUT = 2 * 60 * 60 * 1000;
private static final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
static {
scheduler.scheduleAtFixedRate(WebSocketServer::cleanInactiveSessions, 1, 1, TimeUnit.HOURS);
}
private static void cleanInactiveSessions() {
final var currentTime = System.currentTimeMillis();
final var iterator = userSessions.entrySet().iterator();
while (iterator.hasNext()) {
final var entry = iterator.next();
final var key = entry.getKey();
final var userData = entry.getValue();
final var lastActiveTime = userData.getLastActiveTime();
if (currentTime - lastActiveTime > INACTIVITY_TIMEOUT) {
try {
userData.closeSession();
log.info("closeSession,userId : {}, userType : {}", key.getUserId(), key.getUserType());
} catch (IOException e) {
log.error("closeSession error,userId : {}, userType : {}", key.getUserId(), key.getUserType(), e);
}
iterator.remove();
}
}
}
public static <T> void sendMessageToUser(UserSessionKey key, SocketMessage<T> socketMessage) throws IOException {
sendMessageToUsers(List.of(key), socketMessage);
}
public static <T> void sendMessageToUsers(List<UserSessionKey> keys, SocketMessage<T> socketMessage) throws IOException {
for (final var key : keys) {
final var userData = userSessions.get(key);
if (userData != null && userData.getSession().isOpen()) {
userData.getSession().sendMessage(new TextMessage(socketMessage.toJsonString()));
userData.updateLastActiveTime();
log.info("send message to userId : {}, userType : {}", key.getUserId(), key.getUserType());
} else {
log.info("userId : {}, userType : {} WebSocket connect closed; ", key.getUserId(), key.getUserType());
}
}
}
public static <T> void sendMessageByUserType(UserTypeEnum userType, SocketMessage<T> socketMessage) throws IOException {
sendMessageByUserTypes(List.of(userType), socketMessage);
}
public static <T> void sendMessageByUserTypes(List<UserTypeEnum> userTypes, SocketMessage<T> socketMessage) throws IOException {
final var keys = userSessions.keySet().stream()
.filter(userSessionData -> userTypes.contains(userSessionData.getUserType()))
.collect(Collectors.toList());
sendMessageToUsers(keys, socketMessage);
}
public static void closeConnection(UserSessionKey key) {
final var userData = userSessions.get(key);
if (userData != null) {
try {
userData.closeSession();
userSessions.remove(key);
log.info("WebSocket close; userId = {}, userType = {}", key.getUserId(), key.getUserType());
} catch (IOException e) {
log.error("WebSocket close failed: userId = {}, userType = {}", key.getUserId(), key.getUserType(), e);
}
} else {
log.warn("Not Found userId : {}, userType : {} WebSocket Connect ", key.getUserId(), key.getUserType());
}
}
private UserSessionKey getUserSessionKey(WebSocketSession session) {
final var userId = (String) session.getAttributes().get(CommonConstants.USER_ID);
final var userTypeStr = (String) session.getAttributes().get(CommonConstants.USER_TYPE);
final var userType = UserTypeEnum.fromString(userTypeStr);
if (userId != null && userType != null) {
return UserSessionKey.builder().userId(userId).userType(userType).build();
}
return null;
}
@Override
public void afterConnectionEstablished(WebSocketSession session) {
final var key = getUserSessionKey(session);
if (key != null) {
userSessions.put(key, new UserSessionData(session));
log.info("WebSocket connect success: userId : {}, userType : {}", key.getUserId(), key.getUserType());
} else {
log.warn("WebSocket connect failed,unable to get valid user information");
}
}
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
final var key = getUserSessionKey(session);
if (key != null) {
userSessions.remove(key);
log.info("WebSocket closed: userId : {}, userType : {}", key.getUserId(), key.getUserType());
} else {
log.warn("WebSocket closed failed,unable to get valid user information");
}
}
}
\ No newline at end of file
package com.infoloop.tianting.server.message;
import lombok.Getter;
@Getter
public enum MessageTypeEnum {
ORDER("订单类型消息");
private final String description;
MessageTypeEnum(String description) {
this.description = description;
}
}
package com.infoloop.tianting.server.message;
import lombok.Getter;
@Getter
public enum NoticeTypeEnum {
ORDER_CREATED("下单成功", MessageTypeEnum.ORDER),
ORDER_PAYMENT_SUCCESS("支付成功", MessageTypeEnum.ORDER);
private final String description;
private final MessageTypeEnum messageType;
NoticeTypeEnum(String description, MessageTypeEnum messageType) {
this.description = description;
this.messageType = messageType;
}
}
package com.infoloop.tianting.server.message;
import com.infoloop.tianting.utils.JsonUtil;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.experimental.SuperBuilder;
import java.io.Serializable;
@Data
@SuperBuilder
@NoArgsConstructor
@AllArgsConstructor
public final class SocketMessage<T> implements Serializable {
private static final long serialVersionUID = 1L;
private MessageTypeEnum type;
private NoticeTypeEnum noticeType;
private String subject;
private String message;
private T data;
public String toJsonString() {
return JsonUtil.writeAsJson(this);
}
}
package com.infoloop.tianting.server.message.data;
import lombok.Builder;
import lombok.Data;
@Data
@Builder
public class OrderCreateSuccessData {
private int orderId;
}
package com.infoloop.tianting.server.message.data;
import lombok.Builder;
import lombok.Data;
@Data
@Builder
public class OrderPaymentSuccessData {
private int orderId;
private String orderCode;
private String createTime;
private String payTime;
}
package com.infoloop.tianting.server.session;
import lombok.Getter;
import org.springframework.web.socket.CloseStatus;
import org.springframework.web.socket.WebSocketSession;
import java.io.IOException;
import java.util.concurrent.atomic.AtomicLong;
@Getter
public class UserSessionData {
private final WebSocketSession session;
private final AtomicLong lastActiveTime;
public UserSessionData(WebSocketSession session) {
this.session = session;
this.lastActiveTime = new AtomicLong(System.currentTimeMillis());
}
public long getLastActiveTime() {
return lastActiveTime.get();
}
public void updateLastActiveTime() {
lastActiveTime.set(System.currentTimeMillis());
}
public void closeSession() throws IOException {
if (session != null && session.isOpen()) {
session.close(CloseStatus.GOING_AWAY);
}
}
}
package com.infoloop.tianting.server.session;
import lombok.Builder;
import lombok.Data;
@Data
@Builder
public class UserSessionKey {
private final String userId;
private final UserTypeEnum userType;
}
package com.infoloop.tianting.server.session;
import lombok.Getter;
@Getter
public enum UserTypeEnum {
ENTERPRISE,
CLIENT,
MICRO,
KDS;
public static UserTypeEnum fromString(String value) {
if (value == null || value.isEmpty()) {
return null;
}
try {
return UserTypeEnum.valueOf(value.toUpperCase());
} catch (IllegalArgumentException e) {
return null;
}
}
}
......@@ -42,7 +42,6 @@ import com.infoloop.tianting.clientcustomerorderservice.SkuSellQuantityModificat
import com.infoloop.tianting.clientcustomerorderservice.UpdateClientCustomerOrderRpcRequest;
import com.infoloop.tianting.clientcustomerorderservice.UpdateClientCustomerOrderRpcResponse;
import com.infoloop.tianting.context.LoginContextHolder;
import com.infoloop.tianting.enums.OrderTypeEnum;
import com.infoloop.tianting.exception.ClientEndExceptions;
import com.infoloop.tianting.exception.ErrorCodeEnum;
import com.infoloop.tianting.logic.delay.DelayedQueue;
......@@ -64,17 +63,26 @@ import com.infoloop.tianting.model.dto.OrderDbDTO.OrderDetailOperationDto;
import com.infoloop.tianting.model.dto.OrderDbDTO.QueryOrderByConditionDto;
import com.infoloop.tianting.model.dto.OrderDbDTO.QueryOrderByPaginationDto;
import com.infoloop.tianting.model.dto.OrderDbDTO.UpdateOrderDto;
import com.infoloop.tianting.server.WebSocketServer;
import com.infoloop.tianting.server.message.MessageTypeEnum;
import com.infoloop.tianting.server.message.NoticeTypeEnum;
import com.infoloop.tianting.server.message.SocketMessage;
import com.infoloop.tianting.server.message.data.OrderCreateSuccessData;
import com.infoloop.tianting.server.session.UserTypeEnum;
import com.infoloop.tianting.utils.DateUtil;
import com.infoloop.tianting.utils.JsonUtil;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.redisson.api.RAtomicLong;
import org.apache.commons.lang3.StringUtils;
import org.redisson.api.RedissonClient;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import javax.annotation.Nullable;
import java.io.IOException;
import java.time.Duration;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.util.ArrayList;
import java.util.HashMap;
......@@ -476,19 +484,22 @@ public class OrderServiceRpcClient {
}
public String generatePickupCode() {
LocalDate today = LocalDate.now();
String currentDate = DateUtil.formatDate(today, DateUtil.YYYY_MM_DD);
String redisKey = REDIS_KEY_PREFIX + currentDate;
// 获取分布式原子Long对象,用于原子自增操作
RAtomicLong atomicLong = redissonClient.getAtomicLong(redisKey);
// 原子自增并获取当前值
long counter = atomicLong.incrementAndGet();
final var today = LocalDate.now();
final var currentDate = DateUtil.formatDate(today, DateUtil.YYYY_MM_DD);
final var redisKey = REDIS_KEY_PREFIX + currentDate;
final var atomicLong = redissonClient.getAtomicLong(redisKey);
if (atomicLong.get() == 0) {
LocalDateTime endOfDay = today.atTime(23, 59, 59, 999999999);
long millisUntilEndOfDay = Duration.between(LocalDateTime.now(), endOfDay).toMillis();
atomicLong.expire(Duration.ofMillis(millisUntilEndOfDay));
}
final var counter = atomicLong.incrementAndGet();
return String.format("%04d", (int) counter);
}
public CreateOrderResponse createOrder(CreateOrderDto createOrderDto) {
final var orderDetails = new ArrayList<ClientCustomerOrderDetailCreation>();
for (CreateOrderDetailDto orderDetailDto : createOrderDto.getOrderDetailDtos()) {
for (final var orderDetailDto : createOrderDto.getOrderDetailDtos()) {
final var mealDetail = MealDetail.newBuilder()
.setBasicMaterials(orderDetailDto.getMealDetailJson().getBasicMaterials())
.addAllTastes(orderDetailDto.getMealDetailJson().getTastes())
......@@ -544,8 +555,7 @@ public class OrderServiceRpcClient {
.setOpenId(LoginContextHolder.getOpenId())
.setShouldCreateRoomNo(!createOrderDto.getRoomNo().isEmpty())
.setRoomNo(createOrderDto.getRoomNo())
.setShouldCreateTableCode(createOrderDto.getOrderType().equals(OrderTypeEnum.DINE_IN.getValue()) || !createOrderDto.getTableCode().isEmpty())
.setTableCode(createOrderDto.getTableCode())
.setTableCode(createOrderDto.getTableCode() == null ? "" : createOrderDto.getTableCode())
.setShouldCreateOnlinePay(true)
.setOnlinePay(createOrderDto.getOnlinePay())
.setShouldCreatePayStatus(createOrderDto.getOnlinePay())
......@@ -554,7 +564,7 @@ public class OrderServiceRpcClient {
.setTotalNum(createOrderDto.getTotalNum())
.setPayPrice(createOrderDto.getPayPrice())
.setMealTime(createOrderDto.getMealTime().toLocalDate().atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli())
.setShouldCreateTableCode(createOrderDto.getCloseTimeType().equals(CloseTimeType.IMMEDIATELY))
.setShouldCreateTableCode(createOrderDto.getCloseTimeType().equals(CloseTimeType.IMMEDIATELY) || StringUtils.isNotEmpty(createOrderDto.getTableCode()))
.setShouldCreatePickupCode(createOrderDto.getCloseTimeType().equals(CloseTimeType.IMMEDIATELY))
.setPickupCode(generatePickupCode())
.setShouldCreateRemark(true)
......@@ -593,14 +603,31 @@ public class OrderServiceRpcClient {
throw ClientEndExceptions.BusinessException.build(ErrorCodeEnum.SKUS_STOCK_NO_ENOUGH_TODAY, noEnoughSkus);
}
}
final var orderResponse = orderServiceRpcBlockingStub
.createClientCustomerOrder(CreateClientCustomerOrderRpcRequest.newBuilder()
final var orderResponse = orderServiceRpcBlockingStub.createClientCustomerOrder(CreateClientCustomerOrderRpcRequest.newBuilder()
.setCreation(creation)
.setCreationSource(OrderSourceEnum.forNumber(createOrderDto.getCreationSource()))
.setEnterpriseId(createOrderDto.getEnterpriseId())
.setCreatedBy(createOrderDto.getCreatedBy())
.build());
if (orderResponse.getIsCreated()) {
// 用餐时间为当天,并且非在线支付餐单,发送消息给客户端
if (createOrderDto.getMealTime().toLocalDate().equals(LocalDate.now()) && menuById.getResponse().getOrderRuleJson().getModeOfPayment() == ModeOfPaymentEnum.ONLINE) {
final var orderCreateSuccessData = OrderCreateSuccessData.builder()
.orderId(orderResponse.getId())
.build();
final var data = SocketMessage.<OrderCreateSuccessData>builder()
.type(MessageTypeEnum.ORDER)
.noticeType(NoticeTypeEnum.ORDER_CREATED)
.subject(NoticeTypeEnum.ORDER_CREATED.getDescription())
.message("您下单成功啦,快去查看吧~")
.data(orderCreateSuccessData)
.build();
try {
WebSocketServer.sendMessageByUserTypes(List.of(UserTypeEnum.CLIENT, UserTypeEnum.KDS), data);
} catch (IOException e) {
log.error("Failed to send message to user", e);
}
}
//立即点餐需要扣除库存
if (menuById.getResponse().getOrderRuleJson().getCloseTimeType() == CloseTimeTypeEnum.IMMEDIATELY) {
final var isSuccess = this.addSkusSellQuantities(createOrderDto.getEnterpriseId(), createOrderDto.getStallId(),
......
......@@ -19,6 +19,13 @@ import com.infoloop.tianting.model.dto.PayDTO.PayResponseDto;
import com.infoloop.tianting.model.dto.PayDTO.ThirdPartyPayInformRequestDto;
import com.infoloop.tianting.model.dto.PaymentCallbackDTO;
import com.infoloop.tianting.model.dto.PaymentCallbackResult;
import com.infoloop.tianting.server.WebSocketServer;
import com.infoloop.tianting.server.message.MessageTypeEnum;
import com.infoloop.tianting.server.message.NoticeTypeEnum;
import com.infoloop.tianting.server.message.SocketMessage;
import com.infoloop.tianting.server.message.data.OrderPaymentSuccessData;
import com.infoloop.tianting.server.session.UserSessionKey;
import com.infoloop.tianting.server.session.UserTypeEnum;
import com.infoloop.tianting.utils.DateUtil;
import com.infoloop.tianting.utils.EncryptionUtil;
import lombok.RequiredArgsConstructor;
......@@ -154,6 +161,25 @@ public class PayServiceClient {
.build();
final var updateResponse = orderServiceRpcClient.updateClientCustomerOrderPaymentResult(orderById, build);
log.info("callback updateClientCustomerOrder result: {}", updateResponse.getIsUpdated());
if (updateResponse.getIsUpdated()) {
final var orderPaymentSuccessData = OrderPaymentSuccessData.builder()
.orderId(orderById.getId())
.orderCode(orderById.getOrderCode())
.createTime(DateUtil.formatDate(orderById.getCreatedAt(), DateUtil.Y_M_D_H_M_S))
.payTime(DateUtil.formatDate(orderById.getPayTime(), DateUtil.Y_M_D_H_M_S))
.build();
final var data = SocketMessage.<OrderPaymentSuccessData>builder()
.type(MessageTypeEnum.ORDER)
.noticeType(NoticeTypeEnum.ORDER_PAYMENT_SUCCESS)
.subject(NoticeTypeEnum.ORDER_PAYMENT_SUCCESS.getDescription())
.message("您的订单:" + orderById.getOrderCode() + " 支付成功啦,快去查看吧~")
.data(orderPaymentSuccessData)
.build();
WebSocketServer.sendMessageToUser(UserSessionKey.builder()
.userId(orderById.getOpenId())
.userType(UserTypeEnum.MICRO)
.build(), data);
}
return updateResponse.getIsUpdated();
} catch (JsonProcessingException e) {
log.error("Failed to parse XML response, raw XML: {}", paymentCallbackDTO.getPayResult(), e);
......