commit f8286861cdbc5f7d2ff139b4c5c8de5da5597b2d Author: BrandWang Date: Thu Mar 23 18:05:14 2023 +0800 feat(可回溯): 可回溯 初始化 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..549e00a --- /dev/null +++ b/.gitignore @@ -0,0 +1,33 @@ +HELP.md +target/ +!.mvn/wrapper/maven-wrapper.jar +!**/src/main/**/target/ +!**/src/test/**/target/ + +### STS ### +.apt_generated +.classpath +.factorypath +.project +.settings +.springBeans +.sts4-cache + +### IntelliJ IDEA ### +.idea +*.iws +*.iml +*.ipr + +### NetBeans ### +/nbproject/private/ +/nbbuild/ +/dist/ +/nbdist/ +/.nb-gradle/ +build/ +!**/src/main/**/build/ +!**/src/test/**/build/ + +### VS Code ### +.vscode/ diff --git a/pom.xml b/pom.xml new file mode 100644 index 0000000..5aeb763 --- /dev/null +++ b/pom.xml @@ -0,0 +1,78 @@ + + + 4.0.0 + + org.springframework.boot + spring-boot-starter-parent + 2.3.12.RELEASE + + + com.wabestway.recall + afis-recall + 0.0.1-SNAPSHOT + afis-recall + afis-recall + + 1.8 + + + + org.springframework.boot + spring-boot-starter-actuator + + + org.springframework.boot + spring-boot-starter-data-mongodb-reactive + + + org.springframework.boot + spring-boot-starter-webflux + + + + org.projectlombok + lombok + true + + + org.springframework.boot + spring-boot-starter-test + test + + + io.projectreactor + reactor-test + test + + + com.alibaba + fastjson + 2.0.17 + test + + + org.xerial.snappy + snappy-java + 1.1.8.4 + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + org.projectlombok + lombok + + + + + + + + diff --git a/src/main/java/com/wabestway/recall/AfisRecallApplication.java b/src/main/java/com/wabestway/recall/AfisRecallApplication.java new file mode 100644 index 0000000..9e99165 --- /dev/null +++ b/src/main/java/com/wabestway/recall/AfisRecallApplication.java @@ -0,0 +1,13 @@ +package com.wabestway.recall; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; + +@SpringBootApplication +public class AfisRecallApplication { + + public static void main(String[] args) { + SpringApplication.run(AfisRecallApplication.class, args); + } + +} diff --git a/src/main/java/com/wabestway/recall/model/Order.java b/src/main/java/com/wabestway/recall/model/Order.java new file mode 100644 index 0000000..885fac4 --- /dev/null +++ b/src/main/java/com/wabestway/recall/model/Order.java @@ -0,0 +1,31 @@ +package com.wabestway.recall.model; + +import lombok.Data; +import org.bson.codecs.pojo.annotations.BsonIgnore; +import org.springframework.data.annotation.Id; +import org.springframework.data.mongodb.core.index.Indexed; +import org.springframework.data.mongodb.core.mapping.Document; + +import java.util.List; + +@Data +@Document(collection = "orders") +public class Order { + @Id + private String id; + private String appKey; + private String tenantId; + private String productCode; + private String productName; + @BsonIgnore + private List events; + @Indexed(unique = true) + private String traceId; + @Indexed(unique = true) + private String orderId; + private String complete; + private String archived; + private String fileId; + private String fileUrl; + private long createAt; +} diff --git a/src/main/java/com/wabestway/recall/model/Record.java b/src/main/java/com/wabestway/recall/model/Record.java new file mode 100644 index 0000000..d81df73 --- /dev/null +++ b/src/main/java/com/wabestway/recall/model/Record.java @@ -0,0 +1,29 @@ +package com.wabestway.recall.model; + +import lombok.Data; +import org.bson.codecs.pojo.annotations.BsonIgnore; +import org.springframework.data.annotation.Id; +import org.springframework.data.mongodb.core.index.Indexed; +import org.springframework.data.mongodb.core.mapping.Document; + +import java.util.List; + +@Data +@Document(collection = "records") +public class Record { + @Id + private String id; + private List events; + @Indexed(unique = false) + private String traceId; + @BsonIgnore + private String orderId; + private boolean last; + private String productCode; + private String productName; + private long createAt; + @BsonIgnore + private String appKey; + private String module; + private String content; +} diff --git a/src/main/java/com/wabestway/recall/model/ResObj.java b/src/main/java/com/wabestway/recall/model/ResObj.java new file mode 100644 index 0000000..662d9c3 --- /dev/null +++ b/src/main/java/com/wabestway/recall/model/ResObj.java @@ -0,0 +1,173 @@ +package com.wabestway.recall.model; + +import lombok.Data; + +import java.util.Map; + +@Data +public class ResObj { + private static final String SUCCESS = "success"; + + private static final String FAILURE = "failure"; + + private static final Integer SUCCESS_CODE = 200; + + private static final Integer FAILURE_CODE = 400; + + private Integer code; + + private String message; + + private String _trackId; + + private T data; + + private Map extr; + + private long timestamp; + + private Integer _status; + + private ResObj(Integer code, T data) { + this.code = code; + this.data = data; + this.timestamp = System.currentTimeMillis(); + } + + private ResObj(Integer code, String message, T data) { + this.code = code; + this.data = data; + this.message = message; + this.timestamp = System.currentTimeMillis(); + } + + private ResObj(T data) { + this.data = data; + this.timestamp = System.currentTimeMillis(); + } + + private ResObj(Integer code) { + this.code = code; + this.timestamp = System.currentTimeMillis(); + } + + public boolean isOk() { + return SUCCESS_CODE.equals(this.code); + } + + public T data() { + return (T) this.data; + } + + public Object get(String key) { + if (this.data != null && key != null) { + return ((Map) this.data).get(key); + } + return null; + } + + public ResObj() { + this.timestamp = System.currentTimeMillis(); + } + + public static ResObj ok() { + ResObj o = new ResObj(SUCCESS_CODE); + return o; + } + + public static ResObj ok(Object data) { + ResObj o = new ResObj(SUCCESS_CODE, data); + return o; + } + + public static ResObj fail(String msg) { + ResObj o = new ResObj(FAILURE_CODE, msg, null); + return o; + } + + + public static ResObj fail() { + ResObj o = new ResObj(FAILURE_CODE); + return o; + } + + public static ResObj result(Integer code, Object data) { + ResObj o = new ResObj(code, data); + return o; + } + + public static ResObj result(Integer code, String message, Object data) { + ResObj o = new ResObj(code, message, data); + return o; + } + + public Integer getCode() { + return code; + } + + public void setCode(Integer code) { + this.code = code; + } + + public String getMessage() { + return message; + } + + public void setMessage(String message) { + this.message = message; + } + + public T getData() { + return data; + } + + public ResObj setData(T data) { + this.data = data; + return this; + } + + public Long getTimestamp() { + return timestamp; + } + + public ResObj setTimestamp(Long timestamp) { + this.timestamp = timestamp; + return this; + } + + public Map getExtr() { + return extr; + } + + + public Integer get_status() { + return _status; + } + + public void set_status(Integer _status) { + this._status = _status; + } + + public String get_trackId() { + return _trackId; + } + + public void set_trackId(String _trackId) { + this._trackId = _trackId; + } + + public ResObj trackId(String _trackId) { + this._trackId = _trackId; + return this; + } + + @Override + public String toString() { + if (message == null) { + return "ResObj(" + code + ")"; + } + return "ResObj(" + code + "){" + message + "}"; + } + + +} diff --git a/src/main/java/com/wabestway/recall/repository/OrderRepository.java b/src/main/java/com/wabestway/recall/repository/OrderRepository.java new file mode 100644 index 0000000..45cbe1d --- /dev/null +++ b/src/main/java/com/wabestway/recall/repository/OrderRepository.java @@ -0,0 +1,11 @@ +package com.wabestway.recall.repository; + +import com.wabestway.recall.model.Order; +import org.springframework.data.mongodb.repository.ReactiveMongoRepository; +import reactor.core.publisher.Mono; + +public interface OrderRepository extends ReactiveMongoRepository { + Mono findByOrderId(String orderId); + + Mono findByTraceId(String traceId); +} diff --git a/src/main/java/com/wabestway/recall/repository/RecordRepository.java b/src/main/java/com/wabestway/recall/repository/RecordRepository.java new file mode 100644 index 0000000..43b522c --- /dev/null +++ b/src/main/java/com/wabestway/recall/repository/RecordRepository.java @@ -0,0 +1,11 @@ +package com.wabestway.recall.repository; + + +import com.wabestway.recall.model.Record; +import org.springframework.data.mongodb.repository.ReactiveMongoRepository; +import reactor.core.publisher.Flux; + + +public interface RecordRepository extends ReactiveMongoRepository { + Flux findAllByTraceIdOrderByCreateAtAsc(String traceId); +} diff --git a/src/main/java/com/wabestway/recall/web/OrderController.java b/src/main/java/com/wabestway/recall/web/OrderController.java new file mode 100644 index 0000000..e143f5e --- /dev/null +++ b/src/main/java/com/wabestway/recall/web/OrderController.java @@ -0,0 +1,53 @@ +package com.wabestway.recall.web; + +import com.wabestway.recall.model.Order; +import com.wabestway.recall.model.Record; +import com.wabestway.recall.model.ResObj; +import com.wabestway.recall.repository.OrderRepository; +import com.wabestway.recall.repository.RecordRepository; +import org.springframework.data.domain.Example; +import org.springframework.data.domain.Sort; +import org.springframework.web.bind.annotation.*; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +import java.util.List; + + +@RestController +@RequestMapping("/order") +public class OrderController { + private final RecordRepository recordRepository; + private final OrderRepository orderRepository; + + public OrderController(RecordRepository recordRepository, OrderRepository orderRepository) { + this.recordRepository = recordRepository; + this.orderRepository = orderRepository; + } + + @GetMapping + public Flux getAllOrders() { + return orderRepository.findAll(); + } + + @PostMapping("/list") + public Flux list(@RequestBody Order order) { + Example ep = Example.of(order); + Sort s = Sort.by("createAt"); + return orderRepository.findAll(ep, s); + } + + @PostMapping("/info/{orderId}") + public Mono findByOrderId(@PathVariable String orderId) { + Mono forder = orderRepository.findByOrderId(orderId); + return forder.flatMap(fo -> { + Flux rs = recordRepository.findAllByTraceIdOrderByCreateAtAsc(fo.getTraceId()); + Mono>> es = rs.map(rr -> rr.getEvents()).collectList(); + Mono> events = es.flatMapIterable(lists -> lists).flatMapIterable(list -> list).collectList(); + return events.flatMap(e -> { + fo.setEvents(e); + return Mono.just(ResObj.ok(fo)); + }); + }); + } +} diff --git a/src/main/java/com/wabestway/recall/web/RecallController.java b/src/main/java/com/wabestway/recall/web/RecallController.java new file mode 100644 index 0000000..fbe9572 --- /dev/null +++ b/src/main/java/com/wabestway/recall/web/RecallController.java @@ -0,0 +1,93 @@ +package com.wabestway.recall.web; + +import com.wabestway.recall.model.Order; +import com.wabestway.recall.model.Record; +import com.wabestway.recall.model.ResObj; +import com.wabestway.recall.repository.OrderRepository; +import com.wabestway.recall.repository.RecordRepository; +import org.springframework.web.bind.annotation.*; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +import java.util.UUID; + +@RestController +@RequestMapping("/track") +public class RecallController { + private final RecordRepository recordRepository; + private final OrderRepository orderRepository; + + public RecallController(RecordRepository recordRepository, OrderRepository orderRepository) { + this.recordRepository = recordRepository; + this.orderRepository = orderRepository; + } + + @GetMapping + public Flux getAllUsers() { + return recordRepository.findAll(); + } + + @PostMapping("/save") + public Mono save(@RequestBody Record record) { + if (record.getOrderId() != null && !record.getOrderId().isEmpty()) { + Mono order = orderRepository.findByOrderId(record.getOrderId()); + order.hasElement().subscribe(exit -> { + if (!exit) { + if (record.getTraceId() != null && !record.getTraceId().isEmpty()) { + Mono tcOrder = orderRepository.findByTraceId(record.getTraceId()); + tcOrder.filter(to -> to.getId() != null).flatMap(to -> { + to.setOrderId(record.getOrderId()); + if (record.isLast()) { + to.setComplete("1"); + } + return orderRepository.save(to); + }).subscribe(); + } else { + record.setTraceId(UUID.randomUUID().toString().replace("-", "")); + Order forder = new Order(); + forder.setOrderId(record.getOrderId()); + forder.setTraceId(record.getTraceId()); + forder.setProductCode(record.getProductCode()); + forder.setProductName(record.getProductName()); + forder.setCreateAt(System.currentTimeMillis()); + forder.setAppKey(record.getAppKey()); + forder.setArchived("0"); + forder.setComplete("0"); + if (record.isLast()) { + forder.setComplete("1"); + } + orderRepository.save(forder).subscribe(); + } + } else { + order.subscribe(to -> { + if (record.isLast()) { + to.setComplete("1"); + orderRepository.save(to).subscribe(); + } + }); + } + }); + } else { + if (record.getTraceId() == null || record.getTraceId().isEmpty()) { + record.setTraceId(UUID.randomUUID().toString().replace("-", "")); + Order order = new Order(); + order.setTraceId(record.getTraceId()); + order.setProductCode(record.getProductCode()); + order.setProductName(record.getProductName()); + order.setCreateAt(System.currentTimeMillis()); + order.setAppKey(record.getAppKey()); + order.setArchived("0"); + order.setComplete("0"); + if (record.isLast()) { + order.setComplete("1"); + } + orderRepository.save(order).subscribe(); + } + } + record.setCreateAt(System.currentTimeMillis()); + + return recordRepository.save(record).flatMap(ss -> { + return Mono.just(ResObj.ok(ss)); + }); + } +} diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties new file mode 100644 index 0000000..253dcf4 --- /dev/null +++ b/src/main/resources/application.properties @@ -0,0 +1 @@ +spring.data.mongodb.uri=mongodb://localhost:27017/recall diff --git a/src/test/java/com/wabestway/recall/FFmpeg.java b/src/test/java/com/wabestway/recall/FFmpeg.java new file mode 100644 index 0000000..adb9639 --- /dev/null +++ b/src/test/java/com/wabestway/recall/FFmpeg.java @@ -0,0 +1,12 @@ +package com.wabestway.recall; + + +import org.junit.jupiter.api.Test; + +public class FFmpeg { + @Test + public void TestMp4() { + + + } +}