package com.wabestway.recall.web; import cn.hutool.core.io.FileUtil; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONArray; import com.wabestway.commons.http.ResObj; import com.wabestway.engine.api.dfs.UploadFileDTO; import com.wabestway.engine.api.dfs.UploadRespVO; import com.wabestway.engine.api.feign.DfsStorageFeignClient; import com.wabestway.recall.model.Order; import com.wabestway.recall.model.OrderList; import com.wabestway.recall.model.Record; import com.wabestway.recall.repository.OrderRepository; import com.wabestway.recall.repository.RecordRepository; import lombok.extern.slf4j.Slf4j; import org.springframework.data.domain.*; import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.io.*; import java.nio.file.Files; import java.nio.file.Path; import java.util.ArrayList; import java.util.Comparator; import java.util.List; @Slf4j @RestController @RequestMapping("/order") public class OrderController { private final RecordRepository recordRepository; private final OrderRepository orderRepository; private final DfsStorageFeignClient storageFeignClient; public OrderController(RecordRepository recordRepository, OrderRepository orderRepository, DfsStorageFeignClient storageFeignClient) { this.recordRepository = recordRepository; this.orderRepository = orderRepository; this.storageFeignClient = storageFeignClient; } @GetMapping public Flux getAllOrders() { return orderRepository.findAll(); } @PostMapping("/list") public Mono list(@RequestBody Order order) { OrderList resList = new OrderList(); List list = new ArrayList<>(); Sort s = Sort.by("createAt").descending(); PageRequest pageRequest = PageRequest.of(order.getPage() - 1, order.getPageSize(), s); ExampleMatcher matcher = ExampleMatcher.matching() .withIgnoreNullValues().withIgnorePaths("createAt", "last", "page", "pageSize"); Example ep = Example.of(order, matcher); Flux flux = orderRepository.findAll(ep, s); flux.skip(pageRequest.getOffset()).limitRequest(pageRequest.getPageSize()).subscribe(list::add); resList.setList(list); // Mono> listMono = flux.skip(pageRequest.getOffset()).limitRequest(pageRequest.getPageSize()).collectList(); // return listMono.flatMap(list -> { // resList.setList(list); // return Mono.just(ResObj.ok(resList)); // }); return flux.count().flatMap(count -> { resList.setTotal(count); return Mono.just(ResObj.ok(resList)); }); } @PostMapping("/recallUp") public Mono recallUp(@RequestBody Order order) { Mono queryOrder = orderRepository.findByOrderId(order.getOrderId()); return queryOrder.flatMap(ss -> { ss.setStartDate(order.getStartDate()); ss.setEndDate(order.getEndDate()); ss.setHolderName(order.getHolderName()); ss.setHolderPhone(order.getHolderPhone()); ss.setSupplierName(order.getSupplierName()); orderRepository.save(ss).subscribe(); return Mono.just(ResObj.ok()); }).defaultIfEmpty(ResObj.fail("订单不存在")); } @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)); }); }); } @GetMapping("/ffmpeg") public Mono createMp4() { Order order = new Order(); order.setComplete("1"); ExampleMatcher matcher = ExampleMatcher.matching() .withIgnoreNullValues().withIgnorePaths("createAt", "last", "page", "pageSize"); Example ep = Example.of(order, matcher); Sort s = Sort.by("createAt"); Flux flux = orderRepository.findAll(ep, s); flux.filter(fo -> fo.getFileId() == null).subscribe(fo -> { Mono>> es = recordRepository.findAllByTraceIdOrderByCreateAtAsc(fo.getTraceId()).map(rr -> rr.getEvents()).collectList(); Mono> events = es.flatMapIterable(lists -> lists).flatMapIterable(list -> list).collectList(); events.subscribe(e -> { JSONArray array = JSONArray.parseArray(JSON.toJSONString(e)); try { createVideo(fo, array.toJSONString()); } catch (IOException ex) { ex.printStackTrace(); } }); }); return Mono.just(ResObj.ok()); } void createVideo(Order fo, String jsonData) throws IOException { // 创建临时目录 Path tempDir = null; try { tempDir = Files.createTempDirectory("ts"); } catch (IOException e) { e.printStackTrace(); } System.out.println("Created temporary directory: " + tempDir); try { File tempFile = new File(tempDir.toFile(), "traces.json"); FileWriter fw = new FileWriter(tempFile); BufferedWriter bw = new BufferedWriter(fw); bw.write(jsonData); bw.flush(); bw.close(); fw.close(); String[] command = {"ts-node", "/opt/trace-transform/src/index.ts"}; ProcessBuilder processBuilder = new ProcessBuilder(command); processBuilder.directory(tempDir.toFile()); // 在临时目录下执行命令 Process process = processBuilder.start(); Thread processThread = new Thread(() -> { try { // 读取命令输出 BufferedReader reader = new BufferedReader(new InputStreamReader(process.getInputStream())); String line; while ((line = reader.readLine()) != null) { System.out.println(line); } reader.close(); // 等待命令执行完成 int exitCode = process.waitFor(); System.out.println("Command exited with code " + exitCode); } catch (InterruptedException | IOException e) { e.printStackTrace(); } }); processThread.start(); // Continue with other tasks here... try { processThread.join(); } catch (InterruptedException e) { e.printStackTrace(); } File file = new File(tempDir.toFile(), "video.mp4"); if (file.exists()) { log.info(String.valueOf(file.length())); //byte[] fileBytes = Files.readAllBytes(file.toPath()); byte[] fileBytes = FileUtil.readBytes(file); String titles = fo.getOrderId() + ".mp4"; UploadFileDTO uploadFileDTO = new UploadFileDTO(fileBytes, "RECALL", titles); ResObj uploadResObj = storageFeignClient.uploadBytes(uploadFileDTO); log.info(JSON.toJSONString(uploadResObj)); if (uploadResObj.isOk()) { fo.setFileId(uploadResObj.getData().getFileId()); fo.setFileUrl(uploadResObj.getData().getFileUrl()); orderRepository.save(fo).subscribe(); } } } catch (Exception e) { e.printStackTrace(); log.error(e.getMessage()); } finally { // 删除临时目录 // Files.walk(tempDir) // .sorted(Comparator.reverseOrder()) // .map(Path::toFile) // .forEach(File::delete); } } }