diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTestLogController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTestLogController.java index ebf450a..4f5c7b5 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTestLogController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/aima/AimaTestLogController.java @@ -1,5 +1,10 @@ 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 jakarta.servlet.http.HttpServletResponse; import io.swagger.annotations.Api; @@ -7,6 +12,8 @@ import io.swagger.annotations.ApiOperation; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.beans.factory.annotation.Autowired; 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.core.controller.BaseController; import com.cmvr.common.core.domain.AjaxResult; @@ -23,6 +30,13 @@ public class AimaTestLogController extends BaseController { @Autowired private IAimaTestLogService aimaTestLogService; + @Autowired + private MediaAnalysisClient mediaAnalysisClient; + + @Autowired + @Qualifier("mediaAnalysisTaskExecutor") + private ThreadPoolTaskExecutor mediaAnalysisTaskExecutor; + @ApiOperation("查询测试日志列表") @PreAuthorize("@ss.hasPermi('aima:testlog:list')") @GetMapping("/list") @@ -65,6 +79,84 @@ public class AimaTestLogController extends BaseController { 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("删除测试日志") @PreAuthorize("@ss.hasPermi('aima:testlog:remove')") @Log(title = "测试日志管理", businessType = BusinessType.DELETE) diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java index 2ac53c8..4247d60 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeFlowController.java @@ -2,11 +2,15 @@ package com.cmvr.web.controller.test; import cn.hutool.http.HttpRequest; import cn.hutool.http.HttpResponse; +import org.apache.commons.lang3.StringUtils; import com.alibaba.fastjson2.JSONObject; import com.cmvr.common.annotation.Anonymous; import com.cmvr.common.core.controller.BaseController; import com.cmvr.common.core.domain.AjaxResult; 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.control.FlowControlService; 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.TeTaskExecuteNormalVO; 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.ITeDetectionItemService; 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.ApiOperation; import lombok.RequiredArgsConstructor; +import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; 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.RestController; 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.validation.Valid; import java.util.ArrayList; import java.util.List; +import java.util.UUID; import java.util.concurrent.CancellationException; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; @@ -76,6 +85,9 @@ public class TeFlowController extends BaseController { private final FlowActionExecutorService flowActionExecutorService; private final TiTouchOperateService tiTouchOperateService; private final LLMAiAgentPlatformService llmAiAgentPlatformService; + private final MediaAnalysisClient mediaAnalysisClient; + @Qualifier("mediaAnalysisTaskExecutor") + private final ThreadPoolTaskExecutor mediaAnalysisTaskExecutor; private final List taskInstIdList = new ArrayList<>(); /** @@ -174,6 +186,85 @@ public class TeFlowController extends BaseController { 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("清理终端") @GetMapping("/clear") public AjaxResult clear() { diff --git a/cmvr-iot-admin/src/main/resources/application-aima.yml b/cmvr-iot-admin/src/main/resources/application-aima.yml new file mode 100644 index 0000000..93f4ed2 --- /dev/null +++ b/cmvr-iot-admin/src/main/resources/application-aima.yml @@ -0,0 +1,148 @@ +server: + # 服务器的HTTP端口,默认为8080 + port: 14080 +# 数据源配置 +spring: + datasource: + type: com.alibaba.druid.pool.DruidDataSource + driverClassName: com.mysql.cj.jdbc.Driver + druid: + # 主库数据源 + master: + url: jdbc:mysql://192.168.28.10:3306/${CMVR_DB_NAME:aima}?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8 + username: root + password: cmvr2468 + # 从库数据源 + slave: + # 从数据源开关/默认关闭 + enabled: false + url: + username: + password: + # 初始连接数 + initialSize: 5 + # 最小连接池数量 + minIdle: 10 + # 最大连接池数量 + maxActive: 20 + # 配置获取连接等待超时的时间 + maxWait: 60000 + # 配置连接超时时间 + connectTimeout: 30000 + # 配置网络超时时间 + socketTimeout: 60000 + # 配置间隔多久才进行一次检测,检测需要关闭的空闲连接,单位是毫秒 + timeBetweenEvictionRunsMillis: 60000 + # 配置一个连接在池中最小生存的时间,单位是毫秒 + minEvictableIdleTimeMillis: 300000 + # 配置一个连接在池中最大生存的时间,单位是毫秒 + maxEvictableIdleTimeMillis: 900000 + # 配置检测连接是否有效 + validationQuery: SELECT 1 FROM DUAL + testWhileIdle: true + testOnBorrow: false + testOnReturn: false + webStatFilter: + enabled: true + statViewServlet: + enabled: true + # 设置白名单,不填则允许所有访问 + allow: + url-pattern: /druid/* + # 控制台管理用户名和密码 + login-username: cmvr-iot + login-password: 123456 + filter: + stat: + enabled: true + # 慢SQL记录 + log-slow-sql: true + slow-sql-millis: 1000 + merge-sql: true + wall: + config: + multi-statement-allow: true + # 服务模块 + devtools: + restart: + # 热部署开关 + enabled: true + # redis 配置 + data.redis: + # 地址 + host: 192.168.28.10 + # 端口,默认为6379 + port: 6379 + # 数据库索引 + database: ${CMVR_REDIS_DATABASE:12} + # 密码 + password: cmvr2468 + # 连接超时时间 + timeout: 10s + lettuce: + pool: + # 连接池:中的最小空闲连接 + min-idle: 0 + # 连接池中的最大空闲连接 + max-idle: 8 + # 连接池的最大数据库连接数 + max-active: 8 + # #连接池最大阻塞等待时间(使用负值表示没有限制) + max-wait: -1ms + +# Minio配置 +minio: + url: http://192.168.28.10:9000 + accessKey: AKICMVR + secretKey: wJalrXUtnFEMI + bucketName: cmvr-iot + +# gRPC 客户端设置 +grpc: + client: + grpc-server: # 自定义服务名 + address: 'static://10.148.108.100:50051' # 调用 gRPC 的地址 + # address: 'static://10.148.108.100:50051' # 调用 gRPC 的地址 + enableKeepAlive: true + keepAliveWithoutCalls: true + negotiationType: plaintext # 明文传输 + +api: + app-id: d0epfibvo3em6c4iul40 + app-key: d0eqt1kqek3vg7g4ofi0 + work-flow-url: https://aiagentplatform.cmft.com/api/proxy/api/v1/run_app_workflow + query-result-url: https://aiagentplatform.cmft.com/api/proxy/api/v1/query_run_app_process + bge-vl: http://192.168.1.8:5000/search_similar + generate_advanced_audio: http://192.168.0.102:9003/generate_advanced_audio + # 智能体应用id和秘钥 + # ActionEnum 的枚举作为 key + agents: + INTENT_RECOGNITION: # 意图识别 + app-id: d5dn5ibp9adhq1b34lig + app-key: d6j5bcellh49on5tasvg + TI_TOUCH_COORDINATES: # 获取触控二维坐标 + app-id: d1ebtabnjkflk4gmhikg + app-key: d5thge2cktmipk78h82g + evaluation: http://192.168.0.8:8000/analyze + current-page: http://192.168.0.222:5000/search_similar + +flowise: + tts: 192.168.0.222:8080/tts/ + start: http://192.168.0.108:3000/api/v1/prediction/f99329e9-b33d-437d-90ed-69eaa8a05418 + abort: http://192.168.0.108:3000/api/v1/chatmessage/abort/f99329e9-b33d-437d-90ed-69eaa8a05418/ + query: http://192.168.0.108:3000/api/v1/executions/ + api-key: pg61JW6W_GXyqmqSoURJ4mlSRrCMRzGLBJli_w3BFkg + +media-analysis: + base-url: ${MEDIA_ANALYSIS_BASE_URL:http://192.168.28.10:14080} + api-key: ${MEDIA_ANALYSIS_API_KEY:cmvr-analysis-2026-c65ec181d5ab4a49aef34e6fd1be6704} + connect-timeout-ms: 5000 + read-timeout-ms: 600000 + +# Isolate the resource-center development instance from the currently running robot gateway. +#cmvr: +# resource: +# robot-state: +# enabled: false +# quic: +# enabled: false diff --git a/cmvr-iot-admin/src/main/resources/application-cangan.yml b/cmvr-iot-admin/src/main/resources/application-cangan.yml new file mode 100644 index 0000000..a482a69 --- /dev/null +++ b/cmvr-iot-admin/src/main/resources/application-cangan.yml @@ -0,0 +1,148 @@ +server: + # 服务器的HTTP端口,默认为8080 + port: 14080 +# 数据源配置 +spring: + datasource: + type: com.alibaba.druid.pool.DruidDataSource + driverClassName: com.mysql.cj.jdbc.Driver + druid: + # 主库数据源 + master: + url: jdbc:mysql://192.168.28.10:3306/${CMVR_DB_NAME:changan}?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8 + username: root + password: cmvr2468 + # 从库数据源 + slave: + # 从数据源开关/默认关闭 + enabled: false + url: + username: + password: + # 初始连接数 + initialSize: 5 + # 最小连接池数量 + minIdle: 10 + # 最大连接池数量 + maxActive: 20 + # 配置获取连接等待超时的时间 + maxWait: 60000 + # 配置连接超时时间 + connectTimeout: 30000 + # 配置网络超时时间 + socketTimeout: 60000 + # 配置间隔多久才进行一次检测,检测需要关闭的空闲连接,单位是毫秒 + timeBetweenEvictionRunsMillis: 60000 + # 配置一个连接在池中最小生存的时间,单位是毫秒 + minEvictableIdleTimeMillis: 300000 + # 配置一个连接在池中最大生存的时间,单位是毫秒 + maxEvictableIdleTimeMillis: 900000 + # 配置检测连接是否有效 + validationQuery: SELECT 1 FROM DUAL + testWhileIdle: true + testOnBorrow: false + testOnReturn: false + webStatFilter: + enabled: true + statViewServlet: + enabled: true + # 设置白名单,不填则允许所有访问 + allow: + url-pattern: /druid/* + # 控制台管理用户名和密码 + login-username: cmvr-iot + login-password: 123456 + filter: + stat: + enabled: true + # 慢SQL记录 + log-slow-sql: true + slow-sql-millis: 1000 + merge-sql: true + wall: + config: + multi-statement-allow: true + # 服务模块 + devtools: + restart: + # 热部署开关 + enabled: true + # redis 配置 + data.redis: + # 地址 + host: 192.168.28.10 + # 端口,默认为6379 + port: 6379 + # 数据库索引 + database: ${CMVR_REDIS_DATABASE:11} + # 密码 + password: cmvr2468 + # 连接超时时间 + timeout: 10s + lettuce: + pool: + # 连接池:中的最小空闲连接 + min-idle: 0 + # 连接池中的最大空闲连接 + max-idle: 8 + # 连接池的最大数据库连接数 + max-active: 8 + # #连接池最大阻塞等待时间(使用负值表示没有限制) + max-wait: -1ms + +# Minio配置 +minio: + url: http://192.168.28.10:9000 + accessKey: AKICMVR + secretKey: wJalrXUtnFEMI + bucketName: cmvr-iot + +# gRPC 客户端设置 +grpc: + client: + grpc-server: # 自定义服务名 + address: 'static://10.148.108.100:50051' # 调用 gRPC 的地址 + # address: 'static://10.148.108.100:50051' # 调用 gRPC 的地址 + enableKeepAlive: true + keepAliveWithoutCalls: true + negotiationType: plaintext # 明文传输 + +api: + app-id: d0epfibvo3em6c4iul40 + app-key: d0eqt1kqek3vg7g4ofi0 + work-flow-url: https://aiagentplatform.cmft.com/api/proxy/api/v1/run_app_workflow + query-result-url: https://aiagentplatform.cmft.com/api/proxy/api/v1/query_run_app_process + bge-vl: http://192.168.1.8:5000/search_similar + generate_advanced_audio: http://192.168.0.102:9003/generate_advanced_audio + # 智能体应用id和秘钥 + # ActionEnum 的枚举作为 key + agents: + INTENT_RECOGNITION: # 意图识别 + app-id: d5dn5ibp9adhq1b34lig + app-key: d6j5bcellh49on5tasvg + TI_TOUCH_COORDINATES: # 获取触控二维坐标 + app-id: d1ebtabnjkflk4gmhikg + app-key: d5thge2cktmipk78h82g + evaluation: http://192.168.0.8:8000/analyze + current-page: http://192.168.0.222:5000/search_similar + +flowise: + tts: 192.168.0.222:8080/tts/ + start: http://192.168.0.108:3000/api/v1/prediction/f99329e9-b33d-437d-90ed-69eaa8a05418 + abort: http://192.168.0.108:3000/api/v1/chatmessage/abort/f99329e9-b33d-437d-90ed-69eaa8a05418/ + query: http://192.168.0.108:3000/api/v1/executions/ + api-key: pg61JW6W_GXyqmqSoURJ4mlSRrCMRzGLBJli_w3BFkg + +media-analysis: + base-url: ${MEDIA_ANALYSIS_BASE_URL:http://192.168.28.10:14080} + api-key: ${MEDIA_ANALYSIS_API_KEY:cmvr-analysis-2026-c65ec181d5ab4a49aef34e6fd1be6704} + connect-timeout-ms: 5000 + read-timeout-ms: 600000 + +# Isolate the resource-center development instance from the currently running robot gateway. +cmvr: + resource: + robot-state: + enabled: false + quic: + enabled: false diff --git a/cmvr-iot-admin/src/main/resources/application.yml b/cmvr-iot-admin/src/main/resources/application.yml index c2db5cb..45de4dc 100644 --- a/cmvr-iot-admin/src/main/resources/application.yml +++ b/cmvr-iot-admin/src/main/resources/application.yml @@ -60,7 +60,7 @@ spring: # 国际化资源文件路径 basename: i18n/messages profiles: - active: test + active: aima # 文件上传 servlet: multipart: diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeCameraService.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeCameraService.java index 7a48300..7aaafeb 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeCameraService.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/EdgeCameraService.java @@ -73,6 +73,10 @@ public interface EdgeCameraService { void writeRGBVideoStream(String robotId, String deviceId, OutputStream outputStream); void getDepthImageStream(EdgeCommonVO edgeCommonVO); + + /** Starts the depth camera stream used by the WebSocket channel bridge. */ + StreamObserver getDepthImageStream(String robotId, + String deviceId); void getRGBDImagesStream(EdgeCommonVO edgeCommonVO); diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeCameraServiceImpl.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeCameraServiceImpl.java index 885aa70..86ab433 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeCameraServiceImpl.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeCameraServiceImpl.java @@ -572,7 +572,144 @@ public class EdgeCameraServiceImpl implements EdgeCameraService, EdgeStreamServi @Override 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 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> 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 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 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 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 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 diff --git a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/AudioAnalysisCorrectionRequest.java b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/AudioAnalysisCorrectionRequest.java new file mode 100644 index 0000000..c63914f --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/AudioAnalysisCorrectionRequest.java @@ -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; +} diff --git a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/AudioAnalysisCorrectionSupport.java b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/AudioAnalysisCorrectionSupport.java new file mode 100644 index 0000000..1b6fcb3 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/AudioAnalysisCorrectionSupport.java @@ -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); + } +} diff --git a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClient.java b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClient.java index b554664..ce05c3a 100644 --- a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClient.java +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClient.java @@ -5,5 +5,13 @@ import com.alibaba.fastjson2.JSONObject; public interface MediaAnalysisClient { 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); } diff --git a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClientImpl.java b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClientImpl.java index 80cb3b9..5a88d01 100644 --- a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClientImpl.java +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClientImpl.java @@ -10,6 +10,7 @@ import lombok.RequiredArgsConstructor; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Service; +/** Client for the internal media analysis service. */ @Service @RequiredArgsConstructor public class MediaAnalysisClientImpl implements MediaAnalysisClient { @@ -17,13 +18,21 @@ public class MediaAnalysisClientImpl implements MediaAnalysisClient { private final MediaAnalysisProperties properties; @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()), "/"); if (StringUtils.isBlank(baseUrl)) { - throw new GlobalException("媒体分析服务地址未配置"); + throw new GlobalException(operation + "服务地址未配置"); } - - HttpRequest request = HttpRequest.post(baseUrl + "/api/v1/analysis/run") + HttpRequest request = HttpRequest.post(baseUrl + path) .contentType(ContentType.JSON.toString()) .body(requestBody.toJSONString()) .setConnectionTimeout(properties.getConnectTimeoutMs()) @@ -32,22 +41,21 @@ public class MediaAnalysisClientImpl implements MediaAnalysisClient { if (StringUtils.isNotBlank(apiKey)) { request.bearerAuth(apiKey); } - try (HttpResponse response = request.execute()) { String body = response.body(); if (!response.isOk()) { - String message = extractError(body); - throw new GlobalException("媒体分析服务调用失败({}): {}", response.getStatus(), message); + throw new GlobalException("{}服务调用失败({}): {}", operation, + response.getStatus(), extractError(body)); } JSONObject result = JSON.parseObject(body); if (result == null || !"SUCCEEDED".equalsIgnoreCase(result.getString("status"))) { - throw new GlobalException("媒体分析服务返回了无效结果"); + throw new GlobalException(operation + "服务返回了无效结果"); } return result; } catch (GlobalException exception) { throw exception; } catch (Exception exception) { - throw new GlobalException("媒体分析服务不可用: {}", exception.getMessage()); + throw new GlobalException(operation + "服务不可用: {}", exception.getMessage()); } } diff --git a/cmvr-iot-framework/src/main/java/com/cmvr/framework/websocket/service/MessagePushService.java b/cmvr-iot-framework/src/main/java/com/cmvr/framework/websocket/service/MessagePushService.java index 3d5fc48..8860f33 100644 --- a/cmvr-iot-framework/src/main/java/com/cmvr/framework/websocket/service/MessagePushService.java +++ b/cmvr-iot-framework/src/main/java/com/cmvr/framework/websocket/service/MessagePushService.java @@ -85,7 +85,8 @@ public class MessagePushService { /** 给中途加入频道的会话补发fMP4初始化段。 */ public void replayH264Fmp4Init(String channel, WebSocketSession session) { - if (!channel.contains("/getRGBImageStream/")) { + if (!channel.contains("/getRGBImageStream/") + && !channel.contains("/getDepthImageStream/")) { return; } MediaBootstrap bootstrap = mediaBootstraps.get(channel);