|
@@ -0,0 +1,125 @@
|
|
|
|
|
+package cn.edu.nju.timer_task_service.service.impl;
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+import cn.edu.nju.timer_task_service.data.dao.MessageDAO;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.data.dao.TaskDAO;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.data.entity.Message;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.data.entity.Task;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.service.ChannelService;
|
|
|
|
|
+
|
|
|
|
|
+import cn.edu.nju.timer_task_service.dto.TimerTaskDTO;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.service.TaskService;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.util.ToolKit;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.util.enums.TaskStatus;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.util.exceptions.Asserts;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.util.exceptions.EntityNotAvailableException;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.util.exceptions.ServiceException;
|
|
|
|
|
+import cn.edu.nju.timer_task_service.vo.TaskVO;
|
|
|
|
|
+import lombok.extern.apachecommons.CommonsLog;
|
|
|
|
|
+import org.springframework.beans.BeanUtils;
|
|
|
|
|
+import org.springframework.beans.factory.annotation.Autowired;
|
|
|
|
|
+import org.springframework.beans.factory.annotation.Value;
|
|
|
|
|
+import org.springframework.messaging.support.MessageBuilder;
|
|
|
|
|
+import org.springframework.stereotype.Service;
|
|
|
|
|
+import org.springframework.transaction.annotation.Transactional;
|
|
|
|
|
+
|
|
|
|
|
+import java.time.LocalDateTime;
|
|
|
|
|
+import java.time.ZoneOffset;
|
|
|
|
|
+import java.util.List;
|
|
|
|
|
+
|
|
|
|
|
+/**
|
|
|
|
|
+ * @author fjj
|
|
|
|
|
+ * @date 2019/12/30 10:15 PM
|
|
|
|
|
+ */
|
|
|
|
|
+@Service
|
|
|
|
|
+@CommonsLog
|
|
|
|
|
+public class TaskServiceImpl implements TaskService {
|
|
|
|
|
+ @Autowired
|
|
|
|
|
+ private ChannelService channelService;
|
|
|
|
|
+ @Autowired
|
|
|
|
|
+ private TaskDAO taskDAO;
|
|
|
|
|
+ @Autowired
|
|
|
|
|
+ private MessageDAO messageDAO;
|
|
|
|
|
+
|
|
|
|
|
+ // @Override
|
|
|
|
|
+// @Transactional
|
|
|
|
|
+// public TaskVO create(TimerTaskDTO timerTaskDTO) throws ServiceException {
|
|
|
|
|
+// LocalDateTime now = LocalDateTime.now();
|
|
|
|
|
+// System.out.println(now);
|
|
|
|
|
+// Task task = new Task();
|
|
|
|
|
+// BeanUtils.copyProperties(timerTaskDTO, task);
|
|
|
|
|
+// task.setUuid(ToolKit.randomUUID());
|
|
|
|
|
+// task.setCreatedAt(now);
|
|
|
|
|
+// List<Message> messages = messageDAO.findAllByTask(task);
|
|
|
|
|
+// task = taskDAO.save(task);
|
|
|
|
|
+// Message message = createNewMessage(task);
|
|
|
|
|
+// messages.add(message);
|
|
|
|
|
+// //将消息加入到rabbitmq
|
|
|
|
|
+// long delay = task.getTime().toInstant(ZoneOffset.of("+8")).toEpochMilli() - System.currentTimeMillis();
|
|
|
|
|
+// channelService.output().send(MessageBuilder.withPayload(message.getUuid()).setHeader("x-delay", delay).build());
|
|
|
|
|
+// task.setMessages(messages);
|
|
|
|
|
+// return new TaskVO(taskDAO.save(task));
|
|
|
|
|
+// }
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional
|
|
|
|
|
+ public TaskVO create(TimerTaskDTO timerTaskDTO) throws ServiceException {
|
|
|
|
|
+ LocalDateTime now = LocalDateTime.now();
|
|
|
|
|
+ System.out.println(now);
|
|
|
|
|
+ Task task = new Task();
|
|
|
|
|
+ BeanUtils.copyProperties(timerTaskDTO, task);
|
|
|
|
|
+ task.setUuid(ToolKit.randomUUID());
|
|
|
|
|
+ task.setCreatedAt(now);
|
|
|
|
|
+ Message message = createNewMessage(task);
|
|
|
|
|
+ //将消息加入到rabbitmq
|
|
|
|
|
+ long delay = task.getTime().toInstant(ZoneOffset.of("+8")).toEpochMilli() - System.currentTimeMillis();
|
|
|
|
|
+ channelService.output().send(MessageBuilder.withPayload(message.getUuid()).setHeader("x-delay", delay).build());
|
|
|
|
|
+ task.setMessage(message);
|
|
|
|
|
+ return new TaskVO(taskDAO.save(task));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ @Transactional
|
|
|
|
|
+ public TaskVO update(TimerTaskDTO timerTaskDTO) throws ServiceException {
|
|
|
|
|
+ Task task = taskDAO.findByUuid(timerTaskDTO.getUuid());
|
|
|
|
|
+ Asserts.notNull(task, "该定时任务不存在");
|
|
|
|
|
+ if (task.getStatus() != TaskStatus.NORMAL) {
|
|
|
|
|
+ throw new EntityNotAvailableException("该定时任务已经执行完毕");
|
|
|
|
|
+ }
|
|
|
|
|
+ task.setCallbackUrl(timerTaskDTO.getCallbackUrl());
|
|
|
|
|
+ task.setContent(timerTaskDTO.getContent());
|
|
|
|
|
+ if (!task.getTime().equals(timerTaskDTO.getTime())) {
|
|
|
|
|
+ task.setTime(timerTaskDTO.getTime());
|
|
|
|
|
+ Message oldMessage = messageDAO.findByTask(task);
|
|
|
|
|
+ messageDAO.delete(oldMessage);
|
|
|
|
|
+ // 创建新的message(数据库 and Rabbitmq)
|
|
|
|
|
+ Message newMessage = createNewMessage(task);
|
|
|
|
|
+ //将消息加入到rabbitmq
|
|
|
|
|
+ long delay = task.getTime().toInstant(ZoneOffset.of("+8")).toEpochMilli() - System.currentTimeMillis();
|
|
|
|
|
+ channelService.output().send(MessageBuilder.withPayload(newMessage.getUuid()).setHeader("x-delay", delay).build());
|
|
|
|
|
+ task.setMessage(newMessage);
|
|
|
|
|
+ }
|
|
|
|
|
+ return new TaskVO(taskDAO.save(task));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ public void delete(String uuid) throws ServiceException {
|
|
|
|
|
+ Task task = taskDAO.findByUuid(uuid);
|
|
|
|
|
+ Asserts.notNull(task, "该定时任务不存在");
|
|
|
|
|
+ taskDAO.delete(task);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Override
|
|
|
|
|
+ public TaskVO get(String uuid) throws ServiceException {
|
|
|
|
|
+ Task task = taskDAO.findByUuid(uuid);
|
|
|
|
|
+ Asserts.notNull(task, "该定时任务不存在");
|
|
|
|
|
+ return new TaskVO(task);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private Message createNewMessage(Task task) {
|
|
|
|
|
+ Message message = new Message();
|
|
|
|
|
+ message.setTask(task);
|
|
|
|
|
+ message.setCreatedAt(LocalDateTime.now());
|
|
|
|
|
+ message.setUuid(ToolKit.randomUUID());
|
|
|
|
|
+ return message;
|
|
|
|
|
+ }
|
|
|
|
|
+}
|