소스 검색

refactor: 采集配置下发 + 采集数据上报

zhanghaijun 3 년 전
부모
커밋
a01fecb2d5

+ 31 - 0
src/main/java/com/cec/plugins/controller/CollectController.java

@@ -1,13 +1,21 @@
 package com.cec.plugins.controller;
 
+import cn.hutool.core.bean.BeanUtil;
 import com.cec.plugins.domain.DtoPolicy;
 import com.cec.plugins.domain.ResultData;
+import com.cec.plugins.service.RedisClientSource;
+import com.cec.plugins.service.TaskService;
+import lombok.RequiredArgsConstructor;
+import org.springframework.scheduling.annotation.Async;
 import org.springframework.validation.annotation.Validated;
 import org.springframework.web.bind.annotation.PostMapping;
 import org.springframework.web.bind.annotation.RequestBody;
 import org.springframework.web.bind.annotation.RequestMapping;
 import org.springframework.web.bind.annotation.RestController;
 
+import java.util.ArrayList;
+import java.util.List;
+
 /**
  * @author zhanghaijun
  * @date 2023/8/13 00:27
@@ -18,12 +26,35 @@ import org.springframework.web.bind.annotation.RestController;
  * (4)支持断点续传功能,自动/人工上传功能。
  */
 @RestController
+@RequiredArgsConstructor
 @RequestMapping("/config/api/v1")
 public class CollectController {
 
+    private final RedisClientSource redisClientSource;
+    private final TaskService taskService;
+
     /** 任务状态下发接口 */
     @PostMapping("/sendPolicy")
     public ResultData sendPolicy(@RequestBody @Validated DtoPolicy policy) {
+        reportPolicy(policy);
         return ResultData.success();
     }
+
+    /** 上报采集数据 */
+    @Async
+    public void reportPolicy(DtoPolicy policy) {
+        DtoPolicy.ReportPolicy report = new DtoPolicy.ReportPolicy();
+        List<DtoPolicy.ReportPolicyConfig> configList = new ArrayList<>();
+        BeanUtil.copyProperties(policy.getPolicy(), configList, true);
+        for (DtoPolicy.ReportPolicyConfig reportItem : configList) {
+            reportItem.setScore("80");
+            for (DtoPolicy.ReportPolicyItem item : reportItem.getCollect()) {
+                item.setValue(redisClientSource.getCacheStr(item.getKey()));
+                item.setGroup("运行状态采集项");
+                item.setScore("0");
+            }
+        }
+        report.setData(configList);
+        taskService.reportPlatform("/data/api/v1/collectData", report);
+    }
 }

+ 8 - 7
src/main/java/com/cec/plugins/controller/RedisController.java

@@ -1,5 +1,8 @@
 package com.cec.plugins.controller;
 
+import com.cec.plugins.domain.ResultData;
+import com.cec.plugins.service.RedisClientSource;
+import lombok.RequiredArgsConstructor;
 import org.springframework.web.bind.annotation.GetMapping;
 import org.springframework.web.bind.annotation.PathVariable;
 import org.springframework.web.bind.annotation.RequestMapping;
@@ -8,19 +11,17 @@ import org.springframework.web.bind.annotation.RestController;
 /**
  * @author zhanghaijun
  * @date 2023/8/11 16:04
- * @description [一句话描述该类的功能]
+ * @description 获取redis数据,用于验证数据的正确性
  */
 @RestController
 @RequestMapping("/redis")
+@RequiredArgsConstructor
 public class RedisController {
 
-    @GetMapping("/resource")
-    public String demo() {
-        return "demo";
-    }
+    private final RedisClientSource redisClientSource;
 
     @GetMapping("/resource/{key}")
-    public String resource(@PathVariable String key) {
-        return key;
+    public ResultData resource(@PathVariable String key) {
+        return ResultData.success(redisClientSource.getCacheStr(key));
     }
 }

+ 39 - 2
src/main/java/com/cec/plugins/domain/DtoPolicy.java

@@ -3,8 +3,10 @@ package com.cec.plugins.domain;
 import cn.hutool.json.JSONArray;
 import lombok.Data;
 import lombok.EqualsAndHashCode;
+import lombok.ToString;
 
 import javax.validation.constraints.NotNull;
+import java.util.List;
 
 /**
  * @author zhanghaijun
@@ -13,8 +15,43 @@ import javax.validation.constraints.NotNull;
  */
 @Data
 @EqualsAndHashCode(callSuper = true)
-public class DtoPolicy extends BaseData{
+public class DtoPolicy extends BaseData {
 
     @NotNull(message = "参数policy不能为空")
-    private JSONArray policy;
+    private List<PolicyConfig> policy;
+
+    @Data
+    public static class PolicyConfig {
+        private String deviceId;
+        private String deviceName;
+        private List<PolicyItem> collect;
+    }
+
+    @Data
+    public static class PolicyItem {
+        private String name;
+        private String key;
+        private String value;
+    }
+
+    @Data
+    @EqualsAndHashCode(callSuper = true)
+    public static class ReportPolicy extends BaseData {
+        private List<ReportPolicyConfig> data;
+    }
+
+    @Data
+    public static class ReportPolicyConfig {
+        private String deviceId;
+        private String deviceName;
+        private List<ReportPolicyItem> collect;
+        private String score;
+    }
+
+    @Data
+    @EqualsAndHashCode(callSuper = true)
+    public static class ReportPolicyItem extends PolicyItem {
+        private String group;
+        private String score;
+    }
 }

+ 3 - 5
src/main/java/com/cec/plugins/domain/reported/ReportedData.java

@@ -53,11 +53,11 @@ public class ReportedData extends BaseData {
     }
 
     public interface ReportedType {
-        public TaskEnum getInstanceId();
+        TaskEnum getInstanceId();
 
-        public TaskEnum getStatus();
+        TaskEnum getStatus();
 
-        public TaskEnum getReason();
+        TaskEnum getReason();
     }
 
     public static class ReportedStatus implements ReportedType {
@@ -112,6 +112,4 @@ public class ReportedData extends BaseData {
             return TaskEnum.ExceParamConfigFailReason;
         }
     }
-
-
 }

+ 19 - 27
src/main/java/com/cec/plugins/service/TaskService.java

@@ -53,23 +53,11 @@ public class TaskService {
     /** 注册服务 */
     @Async
     public void register() {
-        ProjectProperties.PlatformModel platform = properties.getPlatform();
-        String url = platform.getHttpUtl("/api/v1/serviceRegistration");
         RegistrationDto body = new RegistrationDto();
         body.setIpAddress(properties.getPlugin().getIpAddress());
         body.setPort(properties.getPlugin().getPort());
         try {
-            Object execute = (new RetryTemplate(platform.getRetries(), platform.getDelayMillis()) {
-                @Override
-                protected Object handle() {
-                    ResultData result = HttpUtils.post(url, body);
-                    if (result.codeSucceed()) {
-                        return result;
-                    }
-                    throw new RetryException(result.getRespMsg());
-                }
-            }).execute();
-            log.info("注册服务成功,返回结果:{}", Convert.toStr(execute));
+            reportPlatform("/api/v1/serviceRegistration", body);
         } catch (Exception e) {
             log.info("注册服务失败,异常信息:{}", e.getMessage());
         }
@@ -105,26 +93,30 @@ public class TaskService {
     public void reported(Class<? extends ReportedData.ReportedType> reportedClass) {
         try {
             Object data = getReportedData(reportedClass).execute();
-            Object result = doReported((ReportedData) data).execute();
-            log.info("{},上报数据成功,返回结果:{}", reportedClass.getName(), Convert.toStr(result));
+            reportPlatform("/config/api/v1/receivetaskstatus", data);
         } catch (Exception e) {
             log.info("{},上报数据失败,异常信息:{}", reportedClass.getName(), e.getMessage());
         }
     }
 
-    private RetryTemplate doReported(ReportedData data) {
-        String url = properties.getPlatform().getHttpUtl("/config/api/v1/receivetaskstatus");
-        ProjectProperties.RetryModel simulation = properties.getSimulation();
-        return new RetryTemplate(simulation.getRetries(), simulation.getDelayMillis()) {
-            @Override
-            protected Object handle() {
-                ResultData result = HttpUtils.post(url, data);
-                if (result.codeSucceed()) {
-                    return result;
+    public void reportPlatform(String apiUrl, Object body) {
+        ProjectProperties.PlatformModel platform = properties.getPlatform();
+        String url = platform.getHttpUtl(apiUrl);
+        try {
+            Object execute = (new RetryTemplate(platform.getRetries(), platform.getDelayMillis()) {
+                @Override
+                protected Object handle() {
+                    ResultData result = HttpUtils.post(url, body);
+                    if (result.codeSucceed()) {
+                        return result;
+                    }
+                    throw new RetryException(result.getRespMsg());
                 }
-                throw new RetryException(result.getRespMsg());
-            }
-        };
+            }).execute();
+            log.info("连接平台成功,请求url:{},返回结果:{}", url, Convert.toStr(execute));
+        } catch (Exception e) {
+            log.info("连接平台失败,请求url:{},异常信息:{}", url, e.getMessage());
+        }
     }
 
     private RetryTemplate getReportedData(Class<? extends ReportedData.ReportedType> reportedClass) {

+ 2 - 2
src/main/resources/logback.xml

@@ -19,8 +19,8 @@
         <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
             <!-- 日志文件名格式 -->
             <fileNamePattern>${log.path}/%d{yyyy-MM-dd}/console.log</fileNamePattern>
-            <!-- 日志最大 1天 -->
-            <maxHistory>1</maxHistory>
+            <!-- 日志最大 7天 -->
+            <maxHistory>7</maxHistory>
         </rollingPolicy>
         <encoder>
             <pattern>${log.pattern}</pattern>