Commit 7cc9acdd by zhangxingmin

push

parent f3787c58
package com.yd.notice.api.controller;
import com.yd.common.result.Result;
import com.yd.notice.api.service.ApiSubscribeRecordService;
import com.yd.notice.feign.request.SubscribeBatchRequest;
import com.yd.notice.feign.request.SubscribeRemainingRequest;
import com.yd.notice.feign.response.SubscribeRemainingResponse;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.*;
import javax.validation.Valid;
/**
* 小程序订阅记录管理
*
* @author zxm
* @since 2026-04-02
*/
@Slf4j
@RestController
@RequestMapping("/subscribe")
@Validated
public class SubscribeRecordController {
@Autowired
private ApiSubscribeRecordService apiSubscribeRecordService;
/**
* 批量保存订阅记录(用户授权后调用)
*/
@PostMapping("/record/batch")
public Result<Void> batchSave(@Valid @RequestBody SubscribeBatchRequest request) {
return apiSubscribeRecordService.batchSave(request);
}
/**
* 查询用户对某个微信模板的剩余有效订阅次数
*/
@GetMapping("/remaining")
public Result<SubscribeRemainingResponse> getRemaining(@Valid SubscribeRemainingRequest request) {
return apiSubscribeRecordService.getRemaining(request);
}
}
\ No newline at end of file
package com.yd.notice.api.service;
import com.yd.common.result.Result;
import com.yd.notice.feign.request.SubscribeBatchRequest;
import com.yd.notice.feign.request.SubscribeRemainingRequest;
import com.yd.notice.feign.response.SubscribeRemainingResponse;
/**
* 小程序订阅记录 API 服务接口
*
* @author zxm
* @since 2026-04-02
*/
public interface ApiSubscribeRecordService {
/**
* 批量保存订阅记录
*/
Result<Void> batchSave(SubscribeBatchRequest request);
/**
* 查询用户对某个微信模板的剩余有效订阅次数
*/
Result<SubscribeRemainingResponse> getRemaining(SubscribeRemainingRequest request);
}
\ No newline at end of file
package com.yd.notice.api.service.impl;
import com.alibaba.fastjson2.JSON;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.yd.common.result.Result;
import com.yd.notice.api.service.ApiSubscribeRecordService;
import com.yd.notice.feign.dto.SubscribeRecordDTO;
import com.yd.notice.feign.request.SubscribeBatchRequest;
import com.yd.notice.feign.request.SubscribeRemainingRequest;
import com.yd.notice.feign.response.SubscribeRemainingResponse;
import com.yd.notice.service.model.NotificationTemplate;
import com.yd.notice.service.model.SubscribeRecord;
import com.yd.notice.service.service.INotificationTemplateService;
import com.yd.notice.service.service.ISubscribeRecordService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.List;
/**
* 小程序订阅记录 API 服务实现类
*
* @author zxm
* @since 2026-04-02
*/
@Slf4j
@Service
public class ApiSubscribeRecordServiceImpl implements ApiSubscribeRecordService {
@Autowired
private ISubscribeRecordService subscribeRecordService;
@Autowired
private INotificationTemplateService templateService;
@Override
public Result<Void> batchSave(SubscribeBatchRequest request) {
log.info("批量保存订阅记录, request={}", JSON.toJSONString(request));
List<SubscribeRecordDTO> dtoList = request.getRecords();
if (dtoList == null || dtoList.isEmpty()) {
return Result.fail("订阅记录列表为空");
}
List<SubscribeRecord> records = new ArrayList<>();
for (SubscribeRecordDTO dto : dtoList) {
// 1. 根据 wxTemplateId 查询模板信息,获取 templateBizId 和 channelBizId
NotificationTemplate template = templateService.getOne(
new LambdaQueryWrapper<NotificationTemplate>()
.eq(NotificationTemplate::getExtraTemplate, dto.getWxTemplateId())
.eq(NotificationTemplate::getStatus, 1)
.orderByDesc(NotificationTemplate::getCreateTime)
.last("LIMIT 1")
);
if (template == null) {
log.error("模板不存在, wxTemplateId={}", dto.getWxTemplateId());
return Result.fail("模板不存在: " + dto.getWxTemplateId());
}
// 2. 构建实体对象
SubscribeRecord record = new SubscribeRecord();
record.setProjectType(dto.getProjectType());
record.setWxTemplateId(dto.getWxTemplateId());
record.setOpenid(dto.getOpenid());
// 从模板表补齐的字段
record.setTemplateBizId(template.getTemplateBizId());
record.setChannelBizId(template.getChannelBizId());
// 系统自动填充字段(默认有效,7天后过期)
record.setStatus(0);
record.setSubscribeTime(LocalDateTime.now());
record.setExpireTime(LocalDateTime.now().plusDays(7));
records.add(record);
}
boolean success = subscribeRecordService.batchSave(records);
if (success) {
log.info("批量保存订阅记录成功, 数量={}", records.size());
return Result.success();
} else {
log.error("批量保存订阅记录失败");
return Result.fail("批量保存失败");
}
}
@Override
public Result<SubscribeRemainingResponse> getRemaining(SubscribeRemainingRequest request) {
log.info("查询剩余次数, request={}", JSON.toJSONString(request));
String openid = request.getOpenid();
String wxTemplateId = request.getWxTemplateId();
// 直接通过 wxTemplateId 统计剩余次数
int count = subscribeRecordService.countRemainingByWxTemplateId(openid, wxTemplateId);
SubscribeRemainingResponse response = new SubscribeRemainingResponse();
response.setOpenid(openid);
response.setWxTemplateId(wxTemplateId);
response.setRemainingCount(count);
log.info("查询剩余次数结果, count={}", count);
return Result.success(response);
}
}
\ No newline at end of file
package com.yd.notice.feign.dto;
import lombok.Data;
import javax.validation.constraints.NotBlank;
/**
* 小程序订阅记录 DTO
*
* @author zxm
* @since 2026-04-02
*/
@Data
public class SubscribeRecordDTO {
/**
* 项目类型(SFP、CFPP等)
*/
@NotBlank(message = "项目类型不能为空")
private String projectType;
/**
* 微信侧实际模板ID(必填,用于查询模板表)
*/
@NotBlank(message = "wxTemplateId模板ID不能为空")
private String wxTemplateId;
/**
* 用户小程序openid(必填)
*/
@NotBlank(message = "用户小程序openid不能为空")
private String openid;
}
\ No newline at end of file
package com.yd.notice.feign.request;
import com.yd.notice.feign.dto.SubscribeRecordDTO;
import lombok.Data;
import javax.validation.Valid;
import javax.validation.constraints.NotEmpty;
import java.util.List;
@Data
public class SubscribeBatchRequest {
/**
* 批量订阅记录列表
*/
@NotEmpty(message = "订阅记录列表不能为空")
@Valid
private List<SubscribeRecordDTO> records;
}
\ No newline at end of file
package com.yd.notice.feign.request;
import lombok.Data;
import javax.validation.constraints.NotBlank;
/**
* 查询剩余订阅次数请求
*
* @author zxm
* @since 2026-04-02
*/
@Data
public class SubscribeRemainingRequest {
@NotBlank(message = "openid不能为空")
private String openid;
@NotBlank(message = "wxTemplateId模板ID不能为空")
private String wxTemplateId;
}
\ No newline at end of file
package com.yd.notice.feign.response;
import lombok.Data;
/**
* 查询剩余订阅次数响应
*
* @author zxm
* @since 2026-04-02
*/
@Data
public class SubscribeRemainingResponse {
/**
* 微信用户openid
*/
private String openid;
/**
* 模板ID
*/
private String wxTemplateId;
/**
* 用户对某个微信模板的剩余有效订阅次数
*/
private Integer remainingCount;
}
\ No newline at end of file
package com.yd.notice.service.dao;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.yd.notice.service.model.SubscribeRecord;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;
import org.apache.ibatis.annotations.Update;
/**
* 订阅记录 Mapper
*
* @author zxm
* @since 2026-04-02
*/
@Mapper
public interface SubscribeRecordMapper extends BaseMapper<SubscribeRecord> {
/**
* 统计剩余有效次数(未过期、未消耗)
*/
@Select("SELECT COUNT(*) FROM subscribe_record " +
"WHERE openid = #{openid} " +
"AND wx_template_id = #{wxTemplateId} " +
"AND status = 0 " +
"AND expire_time > NOW() " +
"AND is_deleted = 0")
int countRemainingByWxTemplateId(@Param("openid") String openid,
@Param("wxTemplateId") String wxTemplateId);
/**
* 消耗一条有效记录(取最早的一条)
*/
@Update("UPDATE subscribe_record SET status = 1, consume_time = NOW() " +
"WHERE id = ( " +
" SELECT id FROM ( " +
" SELECT id FROM subscribe_record " +
" WHERE openid = #{openid} " +
" AND wx_template_id = #{wxTemplateId} " +
" AND status = 0 " +
" AND expire_time > NOW() " +
" AND is_deleted = 0 " +
" ORDER BY subscribe_time ASC LIMIT 1 " +
" ) AS tmp " +
")")
int consumeOneByWxTemplateId(@Param("openid") String openid,
@Param("wxTemplateId") String wxTemplateId);
}
\ No newline at end of file
package com.yd.notice.service.model;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import java.io.Serializable;
import java.time.LocalDateTime;
import lombok.Getter;
import lombok.Setter;
/**
* <p>
* 小程序订阅记录表
* </p>
*
* @author zxm
* @since 2026-08-27
*/
@Getter
@Setter
@TableName("subscribe_record")
public class SubscribeRecord implements Serializable {
private static final long serialVersionUID = 1L;
/**
* 主键ID
*/
@TableId(value = "id", type = IdType.AUTO)
private Long id;
/**
* 项目类型(SFP、CFPP等)
*/
@TableField("project_type")
private String projectType;
/**
* 项目唯一业务ID
*/
@TableField("project_biz_id")
private String projectBizId;
/**
* 渠道业务ID(关联channel_config)
*/
@TableField("channel_biz_id")
private String channelBizId;
/**
* 模板业务ID(关联notification_template)
*/
@TableField("template_biz_id")
private String templateBizId;
/**
* 微信侧实际模板ID(extra_template冗余,便于查询)
*/
@TableField("wx_template_id")
private String wxTemplateId;
/**
* 用户openid
*/
@TableField("openid")
private String openid;
/**
* 用户ID
*/
@TableField("user_id")
private String userId;
/**
* 0-有效(未消耗) 1-已消耗
*/
@TableField("status")
private Integer status;
/**
* 授权时间
*/
@TableField("subscribe_time")
private LocalDateTime subscribeTime;
/**
* 过期时间(订阅时间+7天)
*/
@TableField("expire_time")
private LocalDateTime expireTime;
/**
* 消耗时间
*/
@TableField("consume_time")
private LocalDateTime consumeTime;
/**
* 通用备注
*/
@TableField("remark")
private String remark;
/**
* 删除标识: 0-正常, 1-删除
*/
@TableField("is_deleted")
private Integer isDeleted;
/**
* 创建人ID
*/
@TableField("creator_id")
private String creatorId;
/**
* 更新人ID
*/
@TableField("updater_id")
private String updaterId;
/**
* 创建时间
*/
@TableField("create_time")
private LocalDateTime createTime;
/**
* 更新时间
*/
@TableField("update_time")
private LocalDateTime updateTime;
}
package com.yd.notice.service.request;
import com.yd.notice.service.model.SubscribeRecord;
import lombok.Data;
import javax.validation.constraints.NotEmpty;
import java.util.List;
@Data
public class SubscribeBatchRequest {
@NotEmpty(message = "订阅记录列表不能为空")
private List<SubscribeRecord> records;
}
\ No newline at end of file
......@@ -11,6 +11,7 @@ import com.yd.notice.service.model.NotificationTask;
import com.yd.notice.service.model.NotificationTemplate;
import com.yd.notice.service.service.IChannelConfigService;
import com.yd.notice.service.service.INotificationTemplateService;
import com.yd.notice.service.service.ISubscribeRecordService; // 新增导入
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import me.chanjar.weixin.common.error.WxErrorException;
......@@ -27,6 +28,7 @@ public class MiniprogramMessageSender implements MessageSender {
private final IChannelConfigService channelConfigService;
private final INotificationTemplateService templateService;
private final ISubscribeRecordService subscribeRecordService; // 新增注入
// 不可重试的错误码(永久失败)
private static final String USER_REFUSE_ERROR = "43101"; // 用户拒绝接收
......@@ -60,19 +62,17 @@ public class MiniprogramMessageSender implements MessageSender {
}
// 3. 解析 content_template 并渲染占位符
// task.getContent() 已经是渲染后的完整 JSON 字符串(由 TemplateRenderUtil 提前处理)
// 但为了标准化,我们在这里再次解析,顺便提取 page 跳转路径
String renderedContent = task.getContent();
JSONObject contentJson = JSONObject.parseObject(renderedContent);
log.info("进入小程序消息发送器=>消息内容JSON:{}", JSON.toJSONString(contentJson));
// 提取 page(跳转路径),并从 data 中移除
String page = contentJson.getString("page");
log.info("进入小程序消息发送器=>page:{}", page);
contentJson.remove("page");
log.info("进入小程序消息发送器=>移除page:{}", JSON.toJSONString(contentJson));
// 4. 构建微信订阅消息数据结构(新版本使用 List<MsgData>)
// 4. 构建微信订阅消息数据结构
List<WxMaSubscribeMessage.MsgData> dataList = new ArrayList<>();
for (Map.Entry<String, Object> entry : contentJson.entrySet()) {
String key = entry.getKey();
......@@ -92,25 +92,38 @@ public class MiniprogramMessageSender implements MessageSender {
.toUser(task.getReceiver())
.templateId(wxTemplateId)
.data(dataList)
.page(page) // 如果 page 为空,builder 也会处理
.page(page)
.build();
// log.info("进入小程序消息发送器=>构建消息体WxMaSubscribeMessage:{}", JSON.toJSONString(message));
// 6. 调用微信 SDK 发送
WxMaService wxMaService = buildWxMaService(appid, secret);
wxMaService.getMsgService().sendSubscribeMsg(message);
log.info("微信小程序订阅消息发送成功, taskBizId={}, openid={}", task.getTaskBizId(), task.getReceiver());
// ========== 新增:发送成功后消耗一条订阅记录 ==========
try {
String openid = task.getReceiver();
// 使用从模板中获取的 wxTemplateId
boolean consumed = subscribeRecordService.consumeOneByWxTemplateId(openid, wxTemplateId);
if (consumed) {
log.info("消耗订阅记录成功, openid={}, wxTemplateId={}", openid, wxTemplateId);
} else {
log.warn("发送成功但无可用订阅记录可消耗, openid={}, wxTemplateId={}", openid, wxTemplateId);
}
} catch (Exception e) {
// 消耗失败不影响发送结果,仅记录日志
log.error("消耗订阅记录异常, openid={}, wxTemplateId={}", task.getReceiver(), wxTemplateId, e);
}
return SendResult.success();
} catch (WxErrorException e) {
// 微信 API 抛出的具体异常(含错误码)
Integer errorCode = e.getError().getErrorCode(); // 正确获取错误码
Integer errorCode = e.getError().getErrorCode();
String errCode = errorCode != null ? errorCode.toString() : "UNKNOWN";
String errMsg = e.getMessage();
log.info("微信小程序发送失败, taskBizId={}, errCode={}, errMsg={}", task.getTaskBizId(), errCode, errMsg);
// 处理不可重试错误
if (USER_REFUSE_ERROR.equals(errCode) || TEMPLATE_INVALID.equals(errCode)) {
return SendResult.failNonRetryable(errCode, errMsg);
}
......@@ -133,6 +146,6 @@ public class MiniprogramMessageSender implements MessageSender {
@Override
public String getSupportedChannelType() {
return "miniprogram"; // 必须与 channel_config.channel_type 一致
return "miniprogram";
}
}
\ No newline at end of file
package com.yd.notice.service.service;
import com.baomidou.mybatisplus.extension.service.IService;
import com.yd.notice.service.model.SubscribeRecord;
import java.util.List;
/**
* 订阅记录服务接口
*
* @author zxm
* @since 2026-04-02
*/
public interface ISubscribeRecordService extends IService<SubscribeRecord> {
/**
* 批量保存订阅记录
*/
boolean batchSave(List<SubscribeRecord> records);
/**
* 统计剩余有效次数(通过 wxTemplateId)
*/
int countRemainingByWxTemplateId(String openid, String wxTemplateId);
/**
* 消耗一条有效记录(通过 wxTemplateId)
*/
boolean consumeOneByWxTemplateId(String openid, String wxTemplateId);
}
\ No newline at end of file
package com.yd.notice.service.service.impl;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.yd.notice.service.dao.SubscribeRecordMapper;
import com.yd.notice.service.model.SubscribeRecord;
import com.yd.notice.service.service.ISubscribeRecordService;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.util.StringUtils;
import java.util.List;
/**
* 订阅记录服务实现类
*
* @author zxm
* @since 2026-04-02
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class SubscribeRecordServiceImpl extends ServiceImpl<SubscribeRecordMapper, SubscribeRecord> implements ISubscribeRecordService {
@Override
@Transactional(rollbackFor = Exception.class)
public boolean batchSave(List<SubscribeRecord> records) {
if (records == null || records.isEmpty()) {
return false;
}
return saveBatch(records);
}
@Override
public int countRemainingByWxTemplateId(String openid, String wxTemplateId) {
if (!StringUtils.hasText(openid) || !StringUtils.hasText(wxTemplateId)) {
return 0;
}
return baseMapper.countRemainingByWxTemplateId(openid, wxTemplateId);
}
@Override
@Transactional(rollbackFor = Exception.class)
public boolean consumeOneByWxTemplateId(String openid, String wxTemplateId) {
if (!StringUtils.hasText(openid) || !StringUtils.hasText(wxTemplateId)) {
return false;
}
int affected = baseMapper.consumeOneByWxTemplateId(openid, wxTemplateId);
if (affected > 0) {
log.info("消耗订阅记录成功, openid={}, wxTemplateId={}", openid, wxTemplateId);
} else {
log.warn("无可消耗的有效订阅记录, openid={}, wxTemplateId={}", openid, wxTemplateId);
}
return affected > 0;
}
}
\ No newline at end of file
......@@ -21,7 +21,7 @@ public class MyBatisPlusCodeGenerator {
})
.strategyConfig(builder -> {
builder.addInclude(
"channel_config","notification_record","notification_task","notification_template"
"subscribe_record"
)
.entityBuilder()
.enableLombok()
......
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.yd.notice.service.dao.SubscribeRecordMapper">
</mapper>
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment