Compare commits
2 Commits
67731c69f3
...
699bde3bc6
| Author | SHA1 | Date | |
|---|---|---|---|
| 699bde3bc6 | |||
| 9ebd032d87 |
@ -1,5 +1,10 @@
|
|||||||
package com.cmvr.web.controller.aima;
|
package com.cmvr.web.controller.aima;
|
||||||
|
|
||||||
|
import com.alibaba.fastjson2.JSONObject;
|
||||||
|
import com.cmvr.llm.analysis.AudioAnalysisCorrectionRequest;
|
||||||
|
import com.cmvr.llm.analysis.AudioAnalysisCorrectionSupport;
|
||||||
|
import com.cmvr.llm.analysis.MediaAnalysisClient;
|
||||||
|
import org.apache.commons.lang3.StringUtils;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import jakarta.servlet.http.HttpServletResponse;
|
import jakarta.servlet.http.HttpServletResponse;
|
||||||
import io.swagger.annotations.Api;
|
import io.swagger.annotations.Api;
|
||||||
@ -7,6 +12,8 @@ import io.swagger.annotations.ApiOperation;
|
|||||||
import org.springframework.security.access.prepost.PreAuthorize;
|
import org.springframework.security.access.prepost.PreAuthorize;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.web.bind.annotation.*;
|
import org.springframework.web.bind.annotation.*;
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||||
import com.cmvr.common.annotation.Log;
|
import com.cmvr.common.annotation.Log;
|
||||||
import com.cmvr.common.core.controller.BaseController;
|
import com.cmvr.common.core.controller.BaseController;
|
||||||
import com.cmvr.common.core.domain.AjaxResult;
|
import com.cmvr.common.core.domain.AjaxResult;
|
||||||
@ -23,6 +30,13 @@ public class AimaTestLogController extends BaseController {
|
|||||||
@Autowired
|
@Autowired
|
||||||
private IAimaTestLogService aimaTestLogService;
|
private IAimaTestLogService aimaTestLogService;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
private MediaAnalysisClient mediaAnalysisClient;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Qualifier("mediaAnalysisTaskExecutor")
|
||||||
|
private ThreadPoolTaskExecutor mediaAnalysisTaskExecutor;
|
||||||
|
|
||||||
@ApiOperation("查询测试日志列表")
|
@ApiOperation("查询测试日志列表")
|
||||||
@PreAuthorize("@ss.hasPermi('aima:testlog:list')")
|
@PreAuthorize("@ss.hasPermi('aima:testlog:list')")
|
||||||
@GetMapping("/list")
|
@GetMapping("/list")
|
||||||
@ -65,6 +79,84 @@ public class AimaTestLogController extends BaseController {
|
|||||||
return toAjax(aimaTestLogService.updateAimaTestLog(aimaTestLog));
|
return toAjax(aimaTestLogService.updateAimaTestLog(aimaTestLog));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Correct an Aima audio result and submit the stored audio as a training sample.
|
||||||
|
*/
|
||||||
|
@ApiOperation("爱玛声音分析人工改判并重新训练")
|
||||||
|
@PreAuthorize("@ss.hasPermi('aima:testlog:edit')")
|
||||||
|
@Log(title = "声音分析人工改判", businessType = BusinessType.UPDATE)
|
||||||
|
@PostMapping("/audio-correction/{id}")
|
||||||
|
public AjaxResult correctAudio(@PathVariable("id") String id,
|
||||||
|
@jakarta.validation.Valid @RequestBody AudioAnalysisCorrectionRequest correction) {
|
||||||
|
AimaTestLog log = aimaTestLogService.selectAimaTestLogById(id);
|
||||||
|
if (log == null) {
|
||||||
|
return AjaxResult.error("爱玛日志不存在");
|
||||||
|
}
|
||||||
|
JSONObject output = AudioAnalysisCorrectionSupport.parse(log.getActualResult());
|
||||||
|
String mediaUrl = StringUtils.defaultIfBlank(correction.getMediaUrl(),
|
||||||
|
AudioAnalysisCorrectionSupport.mediaUrl(output, null));
|
||||||
|
String profileCode = StringUtils.defaultIfBlank(correction.getProfileCode(),
|
||||||
|
AudioAnalysisCorrectionSupport.profileCode(output, null));
|
||||||
|
if (StringUtils.isBlank(mediaUrl) || StringUtils.isBlank(profileCode)) {
|
||||||
|
return AjaxResult.error("日志缺少音频地址或分析档案");
|
||||||
|
}
|
||||||
|
String requestId = "aima-audio-correction-" + id + "-" + java.util.UUID.randomUUID();
|
||||||
|
JSONObject trainingRequest = AudioAnalysisCorrectionSupport.trainingRequest(
|
||||||
|
requestId, profileCode, mediaUrl, correction);
|
||||||
|
trainingRequest.put("context", new JSONObject()
|
||||||
|
.fluentPut("source", "AIMA_TEST_LOG")
|
||||||
|
.fluentPut("logId", id)
|
||||||
|
.fluentPut("taskInstanceId", log.getTaskInstanceId())
|
||||||
|
.fluentPut("taskId", log.getTaskId())
|
||||||
|
.fluentPut("testCaseId", log.getTestCaseId()));
|
||||||
|
AudioAnalysisCorrectionSupport.apply(output, correction, "RUNNING", null, null);
|
||||||
|
output.getJSONObject("manualCorrection").put("trainingRequestId", requestId);
|
||||||
|
log.setActualResult(output.toJSONString());
|
||||||
|
log.setTestStatus(Boolean.TRUE.equals(output.getBoolean("passed")) ? "1" : "2");
|
||||||
|
log.setLogContent(StringUtils.defaultIfBlank(log.getLogContent(), "") + "(声音分析已人工改判,训练已提交)");
|
||||||
|
aimaTestLogService.updateAimaTestLog(log);
|
||||||
|
try {
|
||||||
|
mediaAnalysisTaskExecutor.execute(() -> trainAimaAudio(id, requestId, trainingRequest));
|
||||||
|
} catch (RuntimeException exception) {
|
||||||
|
updateAimaTraining(id, requestId, "FAILED", null, exception.getMessage());
|
||||||
|
}
|
||||||
|
return AjaxResult.success(output);
|
||||||
|
}
|
||||||
|
|
||||||
|
private void trainAimaAudio(String logId, String requestId, JSONObject trainingRequest) {
|
||||||
|
try {
|
||||||
|
JSONObject training = mediaAnalysisClient.addAudioReference(trainingRequest);
|
||||||
|
updateAimaTraining(logId, requestId, "SUCCEEDED", training, null);
|
||||||
|
} catch (Exception exception) {
|
||||||
|
updateAimaTraining(logId, requestId, "FAILED", null, exception.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void updateAimaTraining(String logId, String requestId, String status,
|
||||||
|
JSONObject training, String error) {
|
||||||
|
AimaTestLog current = aimaTestLogService.selectAimaTestLogById(logId);
|
||||||
|
if (current == null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
JSONObject output = AudioAnalysisCorrectionSupport.parse(current.getActualResult());
|
||||||
|
JSONObject correction = output.getJSONObject("manualCorrection");
|
||||||
|
if (correction == null || !requestId.equals(correction.getString("trainingRequestId"))) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
correction.put("trainingStatus", status);
|
||||||
|
if (training != null) {
|
||||||
|
correction.put("training", training);
|
||||||
|
}
|
||||||
|
if (StringUtils.isNotBlank(error)) {
|
||||||
|
correction.put("trainingError", error);
|
||||||
|
}
|
||||||
|
current.setActualResult(output.toJSONString());
|
||||||
|
current.setTestStatus(Boolean.TRUE.equals(output.getBoolean("passed")) ? "1" : "2");
|
||||||
|
current.setLogContent(StringUtils.defaultIfBlank(current.getLogContent(), "")
|
||||||
|
+ "(训练" + ("SUCCEEDED".equals(status) ? "完成" : "失败") + ")");
|
||||||
|
aimaTestLogService.updateAimaTestLog(current);
|
||||||
|
}
|
||||||
|
|
||||||
@ApiOperation("删除测试日志")
|
@ApiOperation("删除测试日志")
|
||||||
@PreAuthorize("@ss.hasPermi('aima:testlog:remove')")
|
@PreAuthorize("@ss.hasPermi('aima:testlog:remove')")
|
||||||
@Log(title = "测试日志管理", businessType = BusinessType.DELETE)
|
@Log(title = "测试日志管理", businessType = BusinessType.DELETE)
|
||||||
|
|||||||
@ -2,11 +2,15 @@ package com.cmvr.web.controller.test;
|
|||||||
|
|
||||||
import cn.hutool.http.HttpRequest;
|
import cn.hutool.http.HttpRequest;
|
||||||
import cn.hutool.http.HttpResponse;
|
import cn.hutool.http.HttpResponse;
|
||||||
|
import org.apache.commons.lang3.StringUtils;
|
||||||
import com.alibaba.fastjson2.JSONObject;
|
import com.alibaba.fastjson2.JSONObject;
|
||||||
import com.cmvr.common.annotation.Anonymous;
|
import com.cmvr.common.annotation.Anonymous;
|
||||||
import com.cmvr.common.core.controller.BaseController;
|
import com.cmvr.common.core.controller.BaseController;
|
||||||
import com.cmvr.common.core.domain.AjaxResult;
|
import com.cmvr.common.core.domain.AjaxResult;
|
||||||
import com.cmvr.llm.service.LLMAiAgentPlatformService;
|
import com.cmvr.llm.service.LLMAiAgentPlatformService;
|
||||||
|
import com.cmvr.llm.analysis.AudioAnalysisCorrectionRequest;
|
||||||
|
import com.cmvr.llm.analysis.AudioAnalysisCorrectionSupport;
|
||||||
|
import com.cmvr.llm.analysis.MediaAnalysisClient;
|
||||||
import com.cmvr.test.flow.context.TaskContextManager;
|
import com.cmvr.test.flow.context.TaskContextManager;
|
||||||
import com.cmvr.test.flow.control.FlowControlService;
|
import com.cmvr.test.flow.control.FlowControlService;
|
||||||
import com.cmvr.test.flow.runtime.engine.FlowTaskRuntimeService;
|
import com.cmvr.test.flow.runtime.engine.FlowTaskRuntimeService;
|
||||||
@ -19,6 +23,7 @@ import com.cmvr.test.model.vo.TeFlowVersionActionVO;
|
|||||||
import com.cmvr.test.model.vo.TeFlowViewVO;
|
import com.cmvr.test.model.vo.TeFlowViewVO;
|
||||||
import com.cmvr.test.model.vo.TeTaskExecuteNormalVO;
|
import com.cmvr.test.model.vo.TeTaskExecuteNormalVO;
|
||||||
import com.cmvr.test.model.vo.TeTaskExecuteTrailVO;
|
import com.cmvr.test.model.vo.TeTaskExecuteTrailVO;
|
||||||
|
import com.cmvr.test.model.domain.TeNodeInst;
|
||||||
import com.cmvr.test.service.FlowActionExecutorService;
|
import com.cmvr.test.service.FlowActionExecutorService;
|
||||||
import com.cmvr.test.service.ITeDetectionItemService;
|
import com.cmvr.test.service.ITeDetectionItemService;
|
||||||
import com.cmvr.test.service.ITeDetectionItemVersionService;
|
import com.cmvr.test.service.ITeDetectionItemVersionService;
|
||||||
@ -26,6 +31,7 @@ import com.cmvr.test.service.ITeNodeInstService;
|
|||||||
import io.swagger.annotations.Api;
|
import io.swagger.annotations.Api;
|
||||||
import io.swagger.annotations.ApiOperation;
|
import io.swagger.annotations.ApiOperation;
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
|
import org.springframework.security.access.prepost.PreAuthorize;
|
||||||
import org.springframework.web.bind.annotation.GetMapping;
|
import org.springframework.web.bind.annotation.GetMapping;
|
||||||
import org.springframework.web.bind.annotation.PathVariable;
|
import org.springframework.web.bind.annotation.PathVariable;
|
||||||
import org.springframework.web.bind.annotation.PostMapping;
|
import org.springframework.web.bind.annotation.PostMapping;
|
||||||
@ -34,11 +40,14 @@ import org.springframework.web.bind.annotation.RequestMapping;
|
|||||||
import org.springframework.web.bind.annotation.RequestParam;
|
import org.springframework.web.bind.annotation.RequestParam;
|
||||||
import org.springframework.web.bind.annotation.RestController;
|
import org.springframework.web.bind.annotation.RestController;
|
||||||
import org.springframework.web.bind.annotation.PutMapping;
|
import org.springframework.web.bind.annotation.PutMapping;
|
||||||
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
|
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||||
|
|
||||||
import jakarta.annotation.PreDestroy;
|
import jakarta.annotation.PreDestroy;
|
||||||
import jakarta.validation.Valid;
|
import jakarta.validation.Valid;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
import java.util.UUID;
|
||||||
import java.util.concurrent.CancellationException;
|
import java.util.concurrent.CancellationException;
|
||||||
import java.util.concurrent.ExecutionException;
|
import java.util.concurrent.ExecutionException;
|
||||||
import java.util.concurrent.ExecutorService;
|
import java.util.concurrent.ExecutorService;
|
||||||
@ -76,6 +85,9 @@ public class TeFlowController extends BaseController {
|
|||||||
private final FlowActionExecutorService flowActionExecutorService;
|
private final FlowActionExecutorService flowActionExecutorService;
|
||||||
private final TiTouchOperateService tiTouchOperateService;
|
private final TiTouchOperateService tiTouchOperateService;
|
||||||
private final LLMAiAgentPlatformService llmAiAgentPlatformService;
|
private final LLMAiAgentPlatformService llmAiAgentPlatformService;
|
||||||
|
private final MediaAnalysisClient mediaAnalysisClient;
|
||||||
|
@Qualifier("mediaAnalysisTaskExecutor")
|
||||||
|
private final ThreadPoolTaskExecutor mediaAnalysisTaskExecutor;
|
||||||
private final List<String> taskInstIdList = new ArrayList<>();
|
private final List<String> taskInstIdList = new ArrayList<>();
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@ -174,6 +186,85 @@ public class TeFlowController extends BaseController {
|
|||||||
return AjaxResult.success(nodeInstService.view(flowViewVO));
|
return AjaxResult.success(nodeInstService.view(flowViewVO));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Correct a trial audio result and submit the confirmed sample for retraining.
|
||||||
|
* The media URL and profile are read from the persisted node log by default.
|
||||||
|
*/
|
||||||
|
@ApiOperation("试运行声音分析人工改判并重新训练")
|
||||||
|
@PreAuthorize("@ss.hasPermi('test:config:execute')")
|
||||||
|
@PostMapping("/audio-correction/{nodeInstId}")
|
||||||
|
public AjaxResult correctTrialAudio(@PathVariable Long nodeInstId,
|
||||||
|
@Valid @RequestBody AudioAnalysisCorrectionRequest correction) {
|
||||||
|
TeNodeInst node = nodeInstService.getById(nodeInstId);
|
||||||
|
if (node == null) {
|
||||||
|
return AjaxResult.error("试运行日志不存在");
|
||||||
|
}
|
||||||
|
if (!"AUDIO_EVENT_CLASSIFY".equalsIgnoreCase(node.getAction())) {
|
||||||
|
return AjaxResult.error("该日志不是声音分析节点");
|
||||||
|
}
|
||||||
|
JSONObject output = AudioAnalysisCorrectionSupport.parse(node.getParamsOut());
|
||||||
|
JSONObject input = AudioAnalysisCorrectionSupport.parse(node.getParamsIn());
|
||||||
|
String mediaUrl = StringUtils.defaultIfBlank(correction.getMediaUrl(),
|
||||||
|
AudioAnalysisCorrectionSupport.mediaUrl(output, input));
|
||||||
|
String profileCode = StringUtils.defaultIfBlank(correction.getProfileCode(),
|
||||||
|
AudioAnalysisCorrectionSupport.profileCode(output, input));
|
||||||
|
if (StringUtils.isBlank(mediaUrl) || StringUtils.isBlank(profileCode)) {
|
||||||
|
return AjaxResult.error("日志缺少音频地址或分析档案");
|
||||||
|
}
|
||||||
|
String requestId = "trial-audio-correction-" + nodeInstId + "-" + UUID.randomUUID();
|
||||||
|
JSONObject trainingRequest = AudioAnalysisCorrectionSupport.trainingRequest(
|
||||||
|
requestId, profileCode, mediaUrl, correction);
|
||||||
|
trainingRequest.put("context", new JSONObject()
|
||||||
|
.fluentPut("source", "TRIAL_NODE_LOG")
|
||||||
|
.fluentPut("nodeInstId", nodeInstId)
|
||||||
|
.fluentPut("instId", node.getInstId())
|
||||||
|
.fluentPut("nodeId", node.getNodeId()));
|
||||||
|
AudioAnalysisCorrectionSupport.apply(output, correction, "RUNNING", null, null);
|
||||||
|
output.getJSONObject("manualCorrection").put("trainingRequestId", requestId);
|
||||||
|
node.setParamsOut(output.toJSONString());
|
||||||
|
node.setMessage("声音分析已人工改判,训练任务已提交");
|
||||||
|
nodeInstService.updateById(node);
|
||||||
|
try {
|
||||||
|
mediaAnalysisTaskExecutor.execute(() -> trainTrialAudio(
|
||||||
|
nodeInstId, requestId, trainingRequest));
|
||||||
|
} catch (RuntimeException exception) {
|
||||||
|
updateTrialTraining(nodeInstId, requestId, "FAILED", null, exception.getMessage());
|
||||||
|
}
|
||||||
|
return AjaxResult.success(output);
|
||||||
|
}
|
||||||
|
|
||||||
|
private void trainTrialAudio(Long nodeInstId, String requestId, JSONObject trainingRequest) {
|
||||||
|
try {
|
||||||
|
JSONObject training = mediaAnalysisClient.addAudioReference(trainingRequest);
|
||||||
|
updateTrialTraining(nodeInstId, requestId, "SUCCEEDED", training, null);
|
||||||
|
} catch (Exception exception) {
|
||||||
|
updateTrialTraining(nodeInstId, requestId, "FAILED", null, exception.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void updateTrialTraining(Long nodeInstId, String requestId, String status,
|
||||||
|
JSONObject training, String error) {
|
||||||
|
TeNodeInst current = nodeInstService.getById(nodeInstId);
|
||||||
|
if (current == null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
JSONObject output = AudioAnalysisCorrectionSupport.parse(current.getParamsOut());
|
||||||
|
JSONObject correction = output.getJSONObject("manualCorrection");
|
||||||
|
if (correction == null || !requestId.equals(correction.getString("trainingRequestId"))) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
correction.put("trainingStatus", status);
|
||||||
|
if (training != null) {
|
||||||
|
correction.put("training", training);
|
||||||
|
}
|
||||||
|
if (StringUtils.isNotBlank(error)) {
|
||||||
|
correction.put("trainingError", error);
|
||||||
|
}
|
||||||
|
current.setParamsOut(output.toJSONString());
|
||||||
|
current.setMessage("声音分析已人工改判,训练" + ("SUCCEEDED".equals(status) ? "完成" : "失败"));
|
||||||
|
nodeInstService.updateById(current);
|
||||||
|
}
|
||||||
|
|
||||||
@ApiOperation("清理终端")
|
@ApiOperation("清理终端")
|
||||||
@GetMapping("/clear")
|
@GetMapping("/clear")
|
||||||
public AjaxResult clear() {
|
public AjaxResult clear() {
|
||||||
|
|||||||
@ -60,7 +60,7 @@ spring:
|
|||||||
# 国际化资源文件路径
|
# 国际化资源文件路径
|
||||||
basename: i18n/messages
|
basename: i18n/messages
|
||||||
profiles:
|
profiles:
|
||||||
active: test
|
active: aima
|
||||||
# 文件上传
|
# 文件上传
|
||||||
servlet:
|
servlet:
|
||||||
multipart:
|
multipart:
|
||||||
|
|||||||
@ -73,6 +73,10 @@ public interface EdgeCameraService {
|
|||||||
void writeRGBVideoStream(String robotId, String deviceId, OutputStream outputStream);
|
void writeRGBVideoStream(String robotId, String deviceId, OutputStream outputStream);
|
||||||
|
|
||||||
void getDepthImageStream(EdgeCommonVO edgeCommonVO);
|
void getDepthImageStream(EdgeCommonVO edgeCommonVO);
|
||||||
|
|
||||||
|
/** Starts the depth camera stream used by the WebSocket channel bridge. */
|
||||||
|
StreamObserver<CameraCommand.GetDepthImageStreamCommand.Request> getDepthImageStream(String robotId,
|
||||||
|
String deviceId);
|
||||||
void getRGBDImagesStream(EdgeCommonVO edgeCommonVO);
|
void getRGBDImagesStream(EdgeCommonVO edgeCommonVO);
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@ -572,7 +572,144 @@ public class EdgeCameraServiceImpl implements EdgeCameraService, EdgeStreamServi
|
|||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void getDepthImageStream(EdgeCommonVO edgeCommonVO) {
|
public void getDepthImageStream(EdgeCommonVO edgeCommonVO) {
|
||||||
|
getDepthImageStream(edgeCommonVO.getRobotId(), edgeCommonVO.getDeviceId());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Bridges the robot depth stream to the shared H.264 fragmented-MP4 channel.
|
||||||
|
* The WebSocket subscription dispatcher invokes this overload reflectively.
|
||||||
|
*/
|
||||||
|
@Override
|
||||||
|
public StreamObserver<CameraCommand.GetDepthImageStreamCommand.Request> getDepthImageStream(
|
||||||
|
String robotId, String deviceId) {
|
||||||
|
EdgeCommonVO target = new EdgeCommonVO();
|
||||||
|
target.setRobotId(robotId);
|
||||||
|
target.setDeviceId(deviceId);
|
||||||
|
CameraCommand.CameraState cameraState = status(target);
|
||||||
|
if (cameraState.getIsError()) {
|
||||||
|
throw new GlobalException("Camera is in error state");
|
||||||
|
}
|
||||||
|
if (!cameraState.getIsOpened() || !cameraState.getIsStreaming()) {
|
||||||
|
start(target);
|
||||||
|
cameraState = status(target);
|
||||||
|
}
|
||||||
|
|
||||||
|
String channel = "edgeCameraServiceImpl/getDepthImageStream/" + robotId + "/" + deviceId;
|
||||||
|
String streamKey = robotId + deviceId + "edgeCameraServiceImplgetDepthImageStream";
|
||||||
|
AtomicInteger sequence = new AtomicInteger();
|
||||||
|
AtomicBoolean firstMediaChunk = new AtomicBoolean(true);
|
||||||
|
AtomicReference<StreamObserver<CameraCommand.GetDepthImageStreamCommand.Request>> requestObserverRef =
|
||||||
|
new AtomicReference<>();
|
||||||
|
|
||||||
|
H265Fmp4Transcoder transcoder;
|
||||||
|
try {
|
||||||
|
transcoder = new H265Fmp4Transcoder(
|
||||||
|
robotId + "-" + deviceId + "-depth",
|
||||||
|
cameraState.getFps(),
|
||||||
|
chunk -> {
|
||||||
|
int currentSequence = sequence.incrementAndGet();
|
||||||
|
if (firstMediaChunk.compareAndSet(true, false)) {
|
||||||
|
log.info("Depth camera stream emitted first MP4 chunk, robotId: {}, deviceId: {}, bytes: {}",
|
||||||
|
robotId, deviceId, chunk.length);
|
||||||
|
}
|
||||||
|
messagePushService.pushH264Fmp4ToChannel(channel, chunk, currentSequence);
|
||||||
|
},
|
||||||
|
throwable -> {
|
||||||
|
log.error("Depth camera media transcoding failed, robotId: {}, deviceId: {}",
|
||||||
|
robotId, deviceId, throwable);
|
||||||
|
messagePushService.pushToChannel(channel, "500");
|
||||||
|
StreamObserver<CameraCommand.GetDepthImageStreamCommand.Request> observer =
|
||||||
|
requestObserverRef.get();
|
||||||
|
if (observer != null) {
|
||||||
|
observer.onError(Status.INTERNAL.withDescription("Depth media transcoding failed")
|
||||||
|
.withCause(throwable).asRuntimeException());
|
||||||
|
}
|
||||||
|
});
|
||||||
|
transcoder.start();
|
||||||
|
} catch (IOException exception) {
|
||||||
|
throw new GlobalException("Failed to initialize depth camera transcoder: " + exception.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
|
CameraServiceGrpc.CameraServiceStub stub = grpcServiceManager.getGrpcClient(
|
||||||
|
robotId, CameraServiceGrpc.CameraServiceStub.class);
|
||||||
|
StreamObserver<CameraCommand.GetDepthImageStreamCommand.Feedback> responseObserver =
|
||||||
|
new StreamObserver<>() {
|
||||||
|
@Override
|
||||||
|
public void onNext(CameraCommand.GetDepthImageStreamCommand.Feedback value) {
|
||||||
|
if (!value.getHeader().getSuccess()) {
|
||||||
|
transcoder.close();
|
||||||
|
log.warn("Depth camera stream returned an error, robotId: {}, deviceId: {}, message: {}",
|
||||||
|
robotId, deviceId, value.getHeader().getErrorMessage());
|
||||||
|
messagePushService.pushToChannel(channel, "500");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
CameraCommand.FrameData frame = value.getDepthFrame();
|
||||||
|
if (!value.hasDepthFrame() || frame.getData().isEmpty()) {
|
||||||
|
log.warn("Depth camera stream returned an empty frame, robotId: {}, deviceId: {}",
|
||||||
|
robotId, deviceId);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
transcoder.write(frame.getData().toByteArray(), frame.getIsKeyFrame(), frame.getCodec());
|
||||||
|
} catch (IOException exception) {
|
||||||
|
log.error("Failed to write depth camera media, robotId: {}, deviceId: {}",
|
||||||
|
robotId, deviceId, exception);
|
||||||
|
transcoder.close();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onError(Throwable throwable) {
|
||||||
|
transcoder.close();
|
||||||
|
GrpcStreamManager.removeStream(streamKey);
|
||||||
|
Status status = Status.fromThrowable(throwable);
|
||||||
|
if (status.getCode() != Status.Code.CANCELLED) {
|
||||||
|
log.error("Depth camera gRPC stream ended abnormally, robotId: {}, deviceId: {}",
|
||||||
|
robotId, deviceId, throwable);
|
||||||
|
messagePushService.pushToChannel(channel, "500");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onCompleted() {
|
||||||
|
transcoder.close();
|
||||||
|
GrpcStreamManager.removeStream(streamKey);
|
||||||
|
log.info("Depth camera gRPC stream completed, robotId: {}, deviceId: {}", robotId, deviceId);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
StreamObserver<CameraCommand.GetDepthImageStreamCommand.Request> grpcRequestObserver;
|
||||||
|
try {
|
||||||
|
grpcRequestObserver = stub.getDepthImageStream(responseObserver);
|
||||||
|
requestObserverRef.set(grpcRequestObserver);
|
||||||
|
grpcRequestObserver.onNext(CameraCommand.GetDepthImageStreamCommand.Request.newBuilder()
|
||||||
|
.setHeader(EdgeCommonUtil.buildRequest(deviceId))
|
||||||
|
.setEof(false)
|
||||||
|
.build());
|
||||||
|
} catch (RuntimeException exception) {
|
||||||
|
transcoder.close();
|
||||||
|
throw exception;
|
||||||
|
}
|
||||||
|
|
||||||
|
StreamObserver<CameraCommand.GetDepthImageStreamCommand.Request> observer = grpcRequestObserver;
|
||||||
|
return new StreamObserver<>() {
|
||||||
|
@Override
|
||||||
|
public void onNext(CameraCommand.GetDepthImageStreamCommand.Request value) {
|
||||||
|
observer.onNext(value);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onError(Throwable throwable) {
|
||||||
|
transcoder.close();
|
||||||
|
observer.onError(throwable);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void onCompleted() {
|
||||||
|
transcoder.close();
|
||||||
|
observer.onCompleted();
|
||||||
|
}
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
|
|||||||
@ -0,0 +1,28 @@
|
|||||||
|
package com.cmvr.llm.analysis;
|
||||||
|
|
||||||
|
import io.swagger.annotations.ApiModel;
|
||||||
|
import io.swagger.annotations.ApiModelProperty;
|
||||||
|
import jakarta.validation.constraints.NotBlank;
|
||||||
|
import lombok.Data;
|
||||||
|
|
||||||
|
/** Request used when an operator corrects an audio classification result. */
|
||||||
|
@Data
|
||||||
|
@ApiModel("音频分析人工改判请求")
|
||||||
|
public class AudioAnalysisCorrectionRequest {
|
||||||
|
|
||||||
|
@NotBlank
|
||||||
|
@ApiModelProperty(value = "正确的声音类型编码,例如 POWER_ON、POWER_OFF", required = true)
|
||||||
|
private String label;
|
||||||
|
|
||||||
|
@ApiModelProperty("声音类型显示名称")
|
||||||
|
private String labelName;
|
||||||
|
|
||||||
|
@ApiModelProperty("音频地址;为空时使用日志中保存的地址")
|
||||||
|
private String mediaUrl;
|
||||||
|
|
||||||
|
@ApiModelProperty("分析档案编码;为空时使用日志中的 profileCode")
|
||||||
|
private String profileCode;
|
||||||
|
|
||||||
|
@ApiModelProperty("人工改判备注")
|
||||||
|
private String note;
|
||||||
|
}
|
||||||
@ -0,0 +1,103 @@
|
|||||||
|
package com.cmvr.llm.analysis;
|
||||||
|
|
||||||
|
import com.alibaba.fastjson2.JSONObject;
|
||||||
|
import org.apache.commons.lang3.StringUtils;
|
||||||
|
|
||||||
|
import java.util.Locale;
|
||||||
|
|
||||||
|
/** Shared extraction and result mutation rules for trial and Aima audio logs. */
|
||||||
|
public final class AudioAnalysisCorrectionSupport {
|
||||||
|
|
||||||
|
private AudioAnalysisCorrectionSupport() {
|
||||||
|
}
|
||||||
|
|
||||||
|
public static JSONObject parse(String json) {
|
||||||
|
if (StringUtils.isBlank(json)) {
|
||||||
|
return new JSONObject();
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
JSONObject value = JSONObject.parseObject(json);
|
||||||
|
return value == null ? new JSONObject() : value;
|
||||||
|
} catch (RuntimeException ignored) {
|
||||||
|
return new JSONObject();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public static String mediaUrl(JSONObject output, JSONObject input) {
|
||||||
|
String value = first(output, "mediaUrl", "audioUrl");
|
||||||
|
return StringUtils.defaultIfBlank(value, first(input, "mediaUrl", "audioUrl"));
|
||||||
|
}
|
||||||
|
|
||||||
|
public static String profileCode(JSONObject output, JSONObject input) {
|
||||||
|
String value = first(output, "profileCode");
|
||||||
|
return StringUtils.defaultIfBlank(value, first(input, "profileCode"));
|
||||||
|
}
|
||||||
|
|
||||||
|
public static String first(JSONObject value, String... keys) {
|
||||||
|
if (value == null) {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
for (String key : keys) {
|
||||||
|
String candidate = StringUtils.trimToNull(value.getString(key));
|
||||||
|
if (candidate != null) {
|
||||||
|
return candidate;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
public static JSONObject trainingRequest(String requestId, String profileCode,
|
||||||
|
String mediaUrl, AudioAnalysisCorrectionRequest correction) {
|
||||||
|
JSONObject request = new JSONObject();
|
||||||
|
request.put("requestId", requestId);
|
||||||
|
request.put("profileCode", profileCode);
|
||||||
|
request.put("mediaUrl", mediaUrl);
|
||||||
|
request.put("label", correction.getLabel().trim().toUpperCase(Locale.ROOT));
|
||||||
|
if (StringUtils.isNotBlank(correction.getLabelName())) {
|
||||||
|
request.put("labelName", correction.getLabelName().trim());
|
||||||
|
}
|
||||||
|
return request;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Apply the operator's ground truth while retaining the model response for audit. */
|
||||||
|
public static void apply(JSONObject output, AudioAnalysisCorrectionRequest correction,
|
||||||
|
String trainingStatus, JSONObject trainingResult, String trainingError) {
|
||||||
|
JSONObject modelResult = JSONObject.parseObject(output.toJSONString());
|
||||||
|
modelResult.remove("manualCorrection");
|
||||||
|
String label = correction.getLabel().trim().toUpperCase(Locale.ROOT);
|
||||||
|
String labelName = StringUtils.defaultIfBlank(correction.getLabelName(), label);
|
||||||
|
String expectedLabel = StringUtils.trimToNull(output.getString("expectedLabel"));
|
||||||
|
String expectedLabelName = StringUtils.trimToNull(output.getString("expectedLabelName"));
|
||||||
|
boolean passed = expectedLabel == null
|
||||||
|
|| "ANY_VEHICLE_SOUND".equalsIgnoreCase(expectedLabel)
|
||||||
|
? !"HUMAN_SPEECH".equals(label)
|
||||||
|
: expectedLabel.equalsIgnoreCase(label)
|
||||||
|
|| (expectedLabelName != null && expectedLabelName.equalsIgnoreCase(labelName));
|
||||||
|
output.put("label", label);
|
||||||
|
output.put("labelName", labelName);
|
||||||
|
output.put("matched", true);
|
||||||
|
output.put("vehicleSoundDetected", !"HUMAN_SPEECH".equals(label));
|
||||||
|
output.put("passed", passed);
|
||||||
|
output.put("decisionReason", "人工改判");
|
||||||
|
output.put("decisionSource", "MANUAL_REVIEW");
|
||||||
|
output.put("modelResult", modelResult);
|
||||||
|
JSONObject audit = new JSONObject();
|
||||||
|
audit.put("corrected", true);
|
||||||
|
audit.put("label", label);
|
||||||
|
audit.put("labelName", labelName);
|
||||||
|
audit.put("passed", passed);
|
||||||
|
if (expectedLabel != null) {
|
||||||
|
audit.put("expectedLabel", expectedLabel);
|
||||||
|
}
|
||||||
|
audit.put("note", correction.getNote());
|
||||||
|
audit.put("correctedAt", System.currentTimeMillis());
|
||||||
|
audit.put("trainingStatus", trainingStatus);
|
||||||
|
if (trainingResult != null) {
|
||||||
|
audit.put("training", trainingResult);
|
||||||
|
}
|
||||||
|
if (StringUtils.isNotBlank(trainingError)) {
|
||||||
|
audit.put("trainingError", trainingError);
|
||||||
|
}
|
||||||
|
output.put("manualCorrection", audit);
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -5,5 +5,13 @@ import com.alibaba.fastjson2.JSONObject;
|
|||||||
public interface MediaAnalysisClient {
|
public interface MediaAnalysisClient {
|
||||||
|
|
||||||
JSONObject analyze(JSONObject request);
|
JSONObject analyze(JSONObject request);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Submit a human-confirmed audio sample and rebuild the profile.
|
||||||
|
*
|
||||||
|
* @param request requestId, profileCode, mediaUrl and label
|
||||||
|
* @return analysis service training response
|
||||||
|
*/
|
||||||
|
JSONObject addAudioReference(JSONObject request);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -10,6 +10,7 @@ import lombok.RequiredArgsConstructor;
|
|||||||
import org.apache.commons.lang3.StringUtils;
|
import org.apache.commons.lang3.StringUtils;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
|
|
||||||
|
/** Client for the internal media analysis service. */
|
||||||
@Service
|
@Service
|
||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class MediaAnalysisClientImpl implements MediaAnalysisClient {
|
public class MediaAnalysisClientImpl implements MediaAnalysisClient {
|
||||||
@ -17,13 +18,21 @@ public class MediaAnalysisClientImpl implements MediaAnalysisClient {
|
|||||||
private final MediaAnalysisProperties properties;
|
private final MediaAnalysisProperties properties;
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public JSONObject analyze(JSONObject requestBody) {
|
public JSONObject analyze(JSONObject request) {
|
||||||
|
return post("/api/v1/analysis/run", request, "媒体分析");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public JSONObject addAudioReference(JSONObject request) {
|
||||||
|
return post("/api/v1/audio/references", request, "音频样本训练");
|
||||||
|
}
|
||||||
|
|
||||||
|
private JSONObject post(String path, JSONObject requestBody, String operation) {
|
||||||
String baseUrl = StringUtils.removeEnd(StringUtils.trim(properties.getBaseUrl()), "/");
|
String baseUrl = StringUtils.removeEnd(StringUtils.trim(properties.getBaseUrl()), "/");
|
||||||
if (StringUtils.isBlank(baseUrl)) {
|
if (StringUtils.isBlank(baseUrl)) {
|
||||||
throw new GlobalException("媒体分析服务地址未配置");
|
throw new GlobalException(operation + "服务地址未配置");
|
||||||
}
|
}
|
||||||
|
HttpRequest request = HttpRequest.post(baseUrl + path)
|
||||||
HttpRequest request = HttpRequest.post(baseUrl + "/api/v1/analysis/run")
|
|
||||||
.contentType(ContentType.JSON.toString())
|
.contentType(ContentType.JSON.toString())
|
||||||
.body(requestBody.toJSONString())
|
.body(requestBody.toJSONString())
|
||||||
.setConnectionTimeout(properties.getConnectTimeoutMs())
|
.setConnectionTimeout(properties.getConnectTimeoutMs())
|
||||||
@ -32,22 +41,21 @@ public class MediaAnalysisClientImpl implements MediaAnalysisClient {
|
|||||||
if (StringUtils.isNotBlank(apiKey)) {
|
if (StringUtils.isNotBlank(apiKey)) {
|
||||||
request.bearerAuth(apiKey);
|
request.bearerAuth(apiKey);
|
||||||
}
|
}
|
||||||
|
|
||||||
try (HttpResponse response = request.execute()) {
|
try (HttpResponse response = request.execute()) {
|
||||||
String body = response.body();
|
String body = response.body();
|
||||||
if (!response.isOk()) {
|
if (!response.isOk()) {
|
||||||
String message = extractError(body);
|
throw new GlobalException("{}服务调用失败({}): {}", operation,
|
||||||
throw new GlobalException("媒体分析服务调用失败({}): {}", response.getStatus(), message);
|
response.getStatus(), extractError(body));
|
||||||
}
|
}
|
||||||
JSONObject result = JSON.parseObject(body);
|
JSONObject result = JSON.parseObject(body);
|
||||||
if (result == null || !"SUCCEEDED".equalsIgnoreCase(result.getString("status"))) {
|
if (result == null || !"SUCCEEDED".equalsIgnoreCase(result.getString("status"))) {
|
||||||
throw new GlobalException("媒体分析服务返回了无效结果");
|
throw new GlobalException(operation + "服务返回了无效结果");
|
||||||
}
|
}
|
||||||
return result;
|
return result;
|
||||||
} catch (GlobalException exception) {
|
} catch (GlobalException exception) {
|
||||||
throw exception;
|
throw exception;
|
||||||
} catch (Exception exception) {
|
} catch (Exception exception) {
|
||||||
throw new GlobalException("媒体分析服务不可用: {}", exception.getMessage());
|
throw new GlobalException(operation + "服务不可用: {}", exception.getMessage());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -85,7 +85,8 @@ public class MessagePushService {
|
|||||||
|
|
||||||
/** 给中途加入频道的会话补发fMP4初始化段。 */
|
/** 给中途加入频道的会话补发fMP4初始化段。 */
|
||||||
public void replayH264Fmp4Init(String channel, WebSocketSession session) {
|
public void replayH264Fmp4Init(String channel, WebSocketSession session) {
|
||||||
if (!channel.contains("/getRGBImageStream/")) {
|
if (!channel.contains("/getRGBImageStream/")
|
||||||
|
&& !channel.contains("/getDepthImageStream/")) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
MediaBootstrap bootstrap = mediaBootstraps.get(channel);
|
MediaBootstrap bootstrap = mediaBootstraps.get(channel);
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user