chen.yinxiang
Committed by Jiaqi Xia

农信小程序订阅消息记录数据库

......@@ -11,6 +11,7 @@ import brave.propagation.TraceContext;
import brave.sampler.Sampler;
import com.tianting.infoloop.service.grpc.MealOrderGrpcService;
import com.tianting.infoloop.service.grpc.OrderGrpcService;
import com.tianting.infoloop.service.grpc.OrderSubscriptionMessageGrpcService;
import io.grpc.Server;
import io.grpc.ServerInterceptor;
import io.grpc.ServerInterceptors;
......@@ -66,11 +67,13 @@ public class AppConfig implements ApplicationContextAware {
public Server serviceServer(@Value(SERVICE_PORT) final int port,
final ServerInterceptor serverInterceptor,
final OrderGrpcService orderGrpcService,
final MealOrderGrpcService mealOrderGrpcService) {
final MealOrderGrpcService mealOrderGrpcService,
final OrderSubscriptionMessageGrpcService orderSubscriptionMessageGrpcService) {
return NettyServerBuilder
.forPort(port)
.addService(ServerInterceptors.intercept(orderGrpcService, serverInterceptor))
.addService(ServerInterceptors.intercept(mealOrderGrpcService, serverInterceptor))
.addService(ServerInterceptors.intercept(orderSubscriptionMessageGrpcService, serverInterceptor))
.build();
}
}
......
package com.tianting.infoloop.mapper;
import com.tianting.infoloop.inject.MyBaseMapper;
import com.tianting.infoloop.model.db.OrderSubscriptionMessageDb;
import org.apache.ibatis.annotations.Delete;
import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface OrderSubscriptionMessageMapper extends MyBaseMapper<OrderSubscriptionMessageDb> {
/**
* 批量物理删除(绕过逻辑删除)
*/
@Delete("<script>" +
"DELETE FROM order_subscription_message " +
"WHERE enterpriseId = #{enterpriseId} " +
"AND id IN " +
"<foreach collection='ids' item='id' open='(' separator=',' close=')'>" +
"#{id}" +
"</foreach>" +
"</script>")
int physicalDeleteByIds(@Param("ids") List<Long> ids, @Param("enterpriseId") int enterpriseId);
}
package com.tianting.infoloop.model.db;
import com.baomidou.mybatisplus.annotation.FieldFill;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableLogic;
import com.baomidou.mybatisplus.annotation.TableName;
import com.baomidou.mybatisplus.extension.activerecord.Model;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.experimental.Accessors;
import java.time.LocalDateTime;
/**
* order subscription message DB
* @TableName order_subscription_message
*/
@Data
@Accessors(chain = true)
@EqualsAndHashCode(callSuper = true)
@TableName(value = "order_subscription_message")
public class OrderSubscriptionMessageDb extends Model<OrderSubscriptionMessageDb> {
/**
* 主键ID
*/
@TableId(type = IdType.AUTO)
private Long id;
/**
* 企业ID
*/
private Integer enterpriseId;
/**
* 订餐用户ID(dinerId)
*/
private Integer dinerId;
/**
* 微信用户OpenID
*/
private String openId;
/**
* 订阅消息模板ID
*/
private String templateId;
/**
* 跳转路径(小程序 path)
*/
private String jumpPath;
/**
* 关联的订餐周期开始时间(精确到秒)
*/
private LocalDateTime orderPeriodStartDate;
/**
* 关联的订餐周期结束时间(精确到秒)
*/
private LocalDateTime orderPeriodEndDate;
/**
* 授权时间(创建时间)
*/
@TableField(fill = FieldFill.INSERT)
private LocalDateTime createdAt;
/**
* 更新时间
*/
@TableField(fill = FieldFill.INSERT_UPDATE)
private LocalDateTime updatedAt;
/**
* 授权状态,0-待使用(已授权但未发送),1-已使用(已发送,逻辑删除)
*/
@TableLogic
private Boolean isDeleted;
}
package com.tianting.infoloop.service;
import com.baomidou.mybatisplus.extension.service.IService;
import com.infoloop.tianting.ordersubscriptionmessageservice.CreateOrderSubscriptionMessageRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.GetUserOrderSubscriptionMessageHistoryRpcRequest;
import com.tianting.infoloop.model.db.OrderSubscriptionMessageDb;
import java.util.List;
public interface OrderSubscriptionMessageDbService extends IService<OrderSubscriptionMessageDb> {
OrderSubscriptionMessageDb createOrderSubscriptionMessage(CreateOrderSubscriptionMessageRpcRequest request);
List<OrderSubscriptionMessageDb> getUserOrderSubscriptionMessageHistory(GetUserOrderSubscriptionMessageHistoryRpcRequest request);
/**
* 批量逻辑删除授权记录(表示已使用)
* @param ids 要删除的记录ID列表
* @param enterpriseId 企业ID
* @return 受影响的行数
*/
int batchDeleteByIds(java.util.List<Long> ids, int enterpriseId);
/**
* 批量物理删除授权记录
* @param ids 要删除的记录ID列表
* @param enterpriseId 企业ID
* @return 受影响的行数
*/
int batchPhysicalDeleteByIds(java.util.List<Long> ids, int enterpriseId);
}
package com.tianting.infoloop.service.grpc;
import com.baomidou.dynamic.datasource.annotation.DS;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchDeleteOrderSubscriptionMessagesByIdsRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.CreateOrderSubscriptionMessageRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.CreateOrderSubscriptionMessageRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.GetUserOrderSubscriptionMessageHistoryRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.GetUserOrderSubscriptionMessageHistoryRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.OrderSubscriptionMessageRpcResponse;
import com.infoloop.tianting.ordersubscriptionmessageservice.OrderSubscriptionMessageServiceRpcGrpc;
import com.tianting.infoloop.model.db.OrderSubscriptionMessageDb;
import com.tianting.infoloop.service.OrderSubscriptionMessageDbService;
import io.grpc.Status;
import io.grpc.stub.StreamObserver;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.time.ZoneId;
import java.util.stream.Collectors;
import static com.tianting.infoloop.constants.ConfigConstants.DATA_SOURCE_MASTER;
@Slf4j
@Service
@RequiredArgsConstructor(onConstructor_ = @Autowired)
public class OrderSubscriptionMessageGrpcService extends OrderSubscriptionMessageServiceRpcGrpc.OrderSubscriptionMessageServiceRpcImplBase {
private final OrderSubscriptionMessageDbService orderSubscriptionMessageDbService;
@Override
@DS(DATA_SOURCE_MASTER)
public void createOrderSubscriptionMessage(final CreateOrderSubscriptionMessageRpcRequest request,
final StreamObserver<CreateOrderSubscriptionMessageRpcResponse> responseObserver) {
final var builder = CreateOrderSubscriptionMessageRpcResponse.newBuilder();
try {
final var data = orderSubscriptionMessageDbService.createOrderSubscriptionMessage(request);
builder.setId(data.getId());
} catch (final Exception e) {
log.error("createOrderSubscriptionMessage error;", e);
responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
return;
}
responseObserver.onNext(builder.build());
responseObserver.onCompleted();
}
@Override
@DS(DATA_SOURCE_MASTER)
public void getUserOrderSubscriptionMessageHistory(final GetUserOrderSubscriptionMessageHistoryRpcRequest request,
final StreamObserver<GetUserOrderSubscriptionMessageHistoryRpcResponse> responseObserver) {
final var builder = GetUserOrderSubscriptionMessageHistoryRpcResponse.newBuilder();
try {
final var list = orderSubscriptionMessageDbService.getUserOrderSubscriptionMessageHistory(request);
builder.addAllResponse(list.stream()
.map(this::convertToRpcResponse)
.collect(Collectors.toList()));
} catch (final Exception e) {
log.error("getUserOrderSubscriptionMessageHistory error;", e);
responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
return;
}
responseObserver.onNext(builder.build());
responseObserver.onCompleted();
}
private OrderSubscriptionMessageRpcResponse convertToRpcResponse(OrderSubscriptionMessageDb db) {
OrderSubscriptionMessageRpcResponse.Builder builder = OrderSubscriptionMessageRpcResponse.newBuilder();
builder.setId(db.getId());
if (db.getJumpPath() != null) {
builder.setJumpPath(db.getJumpPath());
}
// 转换LocalDateTime为时间戳(毫秒),与项目中其他时间字段保持一致
if (db.getOrderPeriodStartDate() != null) {
long timestamp = db.getOrderPeriodStartDate().atZone(ZoneId.systemDefault())
.toInstant()
.toEpochMilli();
builder.setOrderPeriodStartDate(timestamp);
}
if (db.getOrderPeriodEndDate() != null) {
long timestamp = db.getOrderPeriodEndDate().atZone(ZoneId.systemDefault())
.toInstant()
.toEpochMilli();
builder.setOrderPeriodEndDate(timestamp);
}
// 转换LocalDateTime为时间戳(毫秒)
if (db.getCreatedAt() != null) {
long timestamp = db.getCreatedAt().atZone(ZoneId.systemDefault())
.toInstant()
.toEpochMilli();
builder.setCreatedAt(timestamp);
}
builder.setIsDeleted(db.getIsDeleted() != null && db.getIsDeleted());
return builder.build();
}
@Override
@DS(DATA_SOURCE_MASTER)
public void batchDeleteOrderSubscriptionMessagesByIds(final BatchDeleteOrderSubscriptionMessagesByIdsRpcRequest request,
final StreamObserver<BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse> responseObserver) {
final var builder = BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse.newBuilder();
try {
final int affectedRows = orderSubscriptionMessageDbService.batchDeleteByIds(
request.getIdsList(),
request.getEnterpriseId());
builder.setAffectedRows(affectedRows);
} catch (final Exception e) {
log.error("batchDeleteOrderSubscriptionMessagesByIds error;", e);
responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
return;
}
responseObserver.onNext(builder.build());
responseObserver.onCompleted();
}
@Override
@DS(DATA_SOURCE_MASTER)
public void batchPhysicalDeleteOrderSubscriptionMessagesByIds(
final BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcRequest request,
final StreamObserver<BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse> responseObserver) {
final var builder = BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse.newBuilder();
try {
final int affectedRows = orderSubscriptionMessageDbService.batchPhysicalDeleteByIds(
request.getIdsList(),
request.getEnterpriseId());
builder.setAffectedRows(affectedRows);
} catch (final Exception e) {
log.error("batchPhysicalDeleteOrderSubscriptionMessagesByIds error;", e);
responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage()).asException());
return;
}
responseObserver.onNext(builder.build());
responseObserver.onCompleted();
}
}
package com.tianting.infoloop.service.impl;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.infoloop.tianting.ordersubscriptionmessageservice.CreateOrderSubscriptionMessageRpcRequest;
import com.infoloop.tianting.ordersubscriptionmessageservice.GetUserOrderSubscriptionMessageHistoryRpcRequest;
import com.tianting.infoloop.mapper.OrderSubscriptionMessageMapper;
import com.tianting.infoloop.model.db.OrderSubscriptionMessageDb;
import com.tianting.infoloop.service.OrderSubscriptionMessageDbService;
import com.tianting.infoloop.utils.exception.ErrorCodeConstants;
import com.tianting.infoloop.utils.exception.ServiceExceptionUtil;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.List;
@Slf4j
@Service
@RequiredArgsConstructor(onConstructor_ = @Autowired)
public class OrderSubscriptionMessageDbServiceImpl extends ServiceImpl<OrderSubscriptionMessageMapper, OrderSubscriptionMessageDb> implements OrderSubscriptionMessageDbService {
@Override
public OrderSubscriptionMessageDb createOrderSubscriptionMessage(CreateOrderSubscriptionMessageRpcRequest request) {
OrderSubscriptionMessageDb orderSubscriptionMessageDb = new OrderSubscriptionMessageDb();
orderSubscriptionMessageDb.setEnterpriseId(request.getEnterpriseId());
if (request.getDinerId() != 0) {
orderSubscriptionMessageDb.setDinerId(request.getDinerId());
}
orderSubscriptionMessageDb.setOpenId(request.getOpenId());
orderSubscriptionMessageDb.setTemplateId(request.getTemplateId());
if (!request.getJumpPath().isEmpty()) {
orderSubscriptionMessageDb.setJumpPath(request.getJumpPath());
}
// 转换时间戳为LocalDateTime(精确到秒),与项目中其他时间字段保持一致
LocalDateTime orderPeriodStartDate = Instant.ofEpochMilli(request.getOrderPeriodStartDate())
.atZone(ZoneId.systemDefault())
.toLocalDateTime();
orderSubscriptionMessageDb.setOrderPeriodStartDate(orderPeriodStartDate);
if (request.getOrderPeriodEndDate() > 0) {
LocalDateTime orderPeriodEndDate = Instant.ofEpochMilli(request.getOrderPeriodEndDate())
.atZone(ZoneId.systemDefault())
.toLocalDateTime();
orderSubscriptionMessageDb.setOrderPeriodEndDate(orderPeriodEndDate);
}
boolean res = this.save(orderSubscriptionMessageDb);
if (!res) {
log.error("创建订阅消息授权记录失败, request: {}", request);
throw ServiceExceptionUtil.exception(ErrorCodeConstants.ORDER_GEN_FAIL);
}
return orderSubscriptionMessageDb;
}
@Override
public List<OrderSubscriptionMessageDb> getUserOrderSubscriptionMessageHistory(GetUserOrderSubscriptionMessageHistoryRpcRequest request) {
LambdaQueryWrapper<OrderSubscriptionMessageDb> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(OrderSubscriptionMessageDb::getEnterpriseId, request.getEnterpriseId());
queryWrapper.eq(OrderSubscriptionMessageDb::getOpenId, request.getOpenId());
// 如果传入了 orderPeriodStartDate,查询大于等于该值的记录
if (request.getOrderPeriodStartDate() > 0) {
LocalDateTime orderPeriodStartDate = Instant.ofEpochMilli(request.getOrderPeriodStartDate())
.atZone(ZoneId.systemDefault())
.toLocalDateTime();
queryWrapper.ge(OrderSubscriptionMessageDb::getOrderPeriodStartDate, orderPeriodStartDate);
}
queryWrapper.orderByDesc(OrderSubscriptionMessageDb::getCreatedAt);
return this.list(queryWrapper);
}
@Override
public int batchDeleteByIds(java.util.List<Long> ids, int enterpriseId) {
if (ids == null || ids.isEmpty()) {
return 0;
}
LambdaQueryWrapper<OrderSubscriptionMessageDb> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(OrderSubscriptionMessageDb::getEnterpriseId, enterpriseId);
queryWrapper.in(OrderSubscriptionMessageDb::getId, ids);
// 逻辑删除(使用 MyBatis-Plus 的 @TableLogic)
return this.remove(queryWrapper) ? ids.size() : 0;
}
@Override
public int batchPhysicalDeleteByIds(java.util.List<Long> ids, int enterpriseId) {
if (ids == null || ids.isEmpty()) {
return 0;
}
// 物理删除(使用 Mapper 中的 SQL,绕过逻辑删除)
return this.baseMapper.physicalDeleteByIds(ids, enterpriseId);
}
}
syntax = "proto3";
option java_multiple_files = true;
option java_package = "com.infoloop.tianting.ordersubscriptionmessageservice";
option java_outer_classname = "OrderSubscriptionMessageServiceProto";
option objc_class_prefix = "OP";
package com.infoloop.tianting.ordersubscriptionmessageservice;
service OrderSubscriptionMessageServiceRpc {
// 创建订阅消息授权
rpc CreateOrderSubscriptionMessage (CreateOrderSubscriptionMessageRpcRequest) returns (CreateOrderSubscriptionMessageRpcResponse) {}
// 查询用户授权历史
rpc GetUserOrderSubscriptionMessageHistory (GetUserOrderSubscriptionMessageHistoryRpcRequest) returns (GetUserOrderSubscriptionMessageHistoryRpcResponse) {}
// 批量逻辑删除授权记录(表示已使用)
rpc BatchDeleteOrderSubscriptionMessagesByIds (BatchDeleteOrderSubscriptionMessagesByIdsRpcRequest) returns (BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse) {}
// 批量物理删除授权记录
rpc BatchPhysicalDeleteOrderSubscriptionMessagesByIds (BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcRequest) returns (BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse) {}
}
// 创建订阅消息授权请求
message CreateOrderSubscriptionMessageRpcRequest {
int32 enterpriseId = 1;
int32 dinerId = 2;
string openId = 3;
string templateId = 4;
int64 orderPeriodStartDate = 5; // Unix 时间戳(毫秒),精确到秒
string jumpPath = 6; // 跳转路径(小程序 path),可选
int64 orderPeriodEndDate = 7; // Unix 时间戳(毫秒),精确到秒,可选
}
// 创建订阅消息授权响应
message CreateOrderSubscriptionMessageRpcResponse {
int64 id = 1;
}
// 查询用户授权历史请求
message GetUserOrderSubscriptionMessageHistoryRpcRequest {
int32 enterpriseId = 1;
string openId = 2;
int64 orderPeriodStartDate = 3; // Unix 时间戳(毫秒),精确到秒,可选,如果传入则查询大于等于该值的记录
}
// 订阅消息授权记录
message OrderSubscriptionMessageRpcResponse {
int64 id = 1;
int64 orderPeriodStartDate = 2; // Unix 时间戳(毫秒),精确到秒
int64 createdAt = 3; // Unix 时间戳(毫秒)
bool isDeleted = 4;
string jumpPath = 5; // 跳转路径(小程序 path),可选
int64 orderPeriodEndDate = 6; // Unix 时间戳(毫秒),精确到秒,可选
}
// 查询用户授权历史响应
message GetUserOrderSubscriptionMessageHistoryRpcResponse {
repeated OrderSubscriptionMessageRpcResponse response = 1;
}
// 批量逻辑删除授权记录请求
message BatchDeleteOrderSubscriptionMessagesByIdsRpcRequest {
int32 enterpriseId = 1;
repeated int64 ids = 2; // 要删除的记录ID列表
}
// 批量逻辑删除授权记录响应
message BatchDeleteOrderSubscriptionMessagesByIdsRpcResponse {
int32 affectedRows = 1; // 受影响的行数
}
// 批量物理删除授权记录请求
message BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcRequest {
int32 enterpriseId = 1;
repeated int64 ids = 2; // 要删除的记录ID列表
}
// 批量物理删除授权记录响应
message BatchPhysicalDeleteOrderSubscriptionMessagesByIdsRpcResponse {
int32 affectedRows = 1; // 受影响的行数
}