diff --git a/cmvr-iot-admin/src/main/resources/application-dev.yml b/cmvr-iot-admin/src/main/resources/application-dev.yml index 83386b1..3ebc0e5 100644 --- a/cmvr-iot-admin/src/main/resources/application-dev.yml +++ b/cmvr-iot-admin/src/main/resources/application-dev.yml @@ -132,3 +132,9 @@ flowise: 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:} + connect-timeout-ms: 5000 + read-timeout-ms: 600000 diff --git a/cmvr-iot-admin/src/main/resources/application-prod.yml b/cmvr-iot-admin/src/main/resources/application-prod.yml index 83386b1..3ebc0e5 100644 --- a/cmvr-iot-admin/src/main/resources/application-prod.yml +++ b/cmvr-iot-admin/src/main/resources/application-prod.yml @@ -132,3 +132,9 @@ flowise: 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:} + connect-timeout-ms: 5000 + read-timeout-ms: 600000 diff --git a/cmvr-iot-admin/src/main/resources/application-test.yml b/cmvr-iot-admin/src/main/resources/application-test.yml index 83386b1..dc27746 100644 --- a/cmvr-iot-admin/src/main/resources/application-test.yml +++ b/cmvr-iot-admin/src/main/resources/application-test.yml @@ -132,3 +132,9 @@ flowise: 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 diff --git a/cmvr-iot-api/cmvr-iot-llm/pom.xml b/cmvr-iot-api/cmvr-iot-llm/pom.xml index 27b12a9..d0f10f7 100644 --- a/cmvr-iot-api/cmvr-iot-llm/pom.xml +++ b/cmvr-iot-api/cmvr-iot-llm/pom.xml @@ -30,6 +30,12 @@ 1.0.49 + + org.springframework.boot + spring-boot-starter-test + test + + - \ No newline at end of file + 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 new file mode 100644 index 0000000..b554664 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClient.java @@ -0,0 +1,9 @@ +package com.cmvr.llm.analysis; + +import com.alibaba.fastjson2.JSONObject; + +public interface MediaAnalysisClient { + + JSONObject analyze(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 new file mode 100644 index 0000000..80cb3b9 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisClientImpl.java @@ -0,0 +1,79 @@ +package com.cmvr.llm.analysis; + +import cn.hutool.http.ContentType; +import cn.hutool.http.HttpRequest; +import cn.hutool.http.HttpResponse; +import com.alibaba.fastjson2.JSON; +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; +import lombok.RequiredArgsConstructor; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Service; + +@Service +@RequiredArgsConstructor +public class MediaAnalysisClientImpl implements MediaAnalysisClient { + + private final MediaAnalysisProperties properties; + + @Override + public JSONObject analyze(JSONObject requestBody) { + String baseUrl = StringUtils.removeEnd(StringUtils.trim(properties.getBaseUrl()), "/"); + if (StringUtils.isBlank(baseUrl)) { + throw new GlobalException("媒体分析服务地址未配置"); + } + + HttpRequest request = HttpRequest.post(baseUrl + "/api/v1/analysis/run") + .contentType(ContentType.JSON.toString()) + .body(requestBody.toJSONString()) + .setConnectionTimeout(properties.getConnectTimeoutMs()) + .setReadTimeout(properties.getReadTimeoutMs()); + String apiKey = normalizeApiKey(properties.getApiKey()); + 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); + } + JSONObject result = JSON.parseObject(body); + if (result == null || !"SUCCEEDED".equalsIgnoreCase(result.getString("status"))) { + throw new GlobalException("媒体分析服务返回了无效结果"); + } + return result; + } catch (GlobalException exception) { + throw exception; + } catch (Exception exception) { + throw new GlobalException("媒体分析服务不可用: {}", exception.getMessage()); + } + } + + static String normalizeApiKey(String apiKey) { + String normalized = StringUtils.trimToEmpty(apiKey); + if ("Bearer".equalsIgnoreCase(normalized)) { + return ""; + } + if (StringUtils.startsWithIgnoreCase(normalized, "Bearer")) { + String remainder = normalized.substring("Bearer".length()); + if (StringUtils.isNotBlank(remainder) && Character.isWhitespace(remainder.charAt(0))) { + normalized = StringUtils.trim(remainder); + } + } + return normalized; + } + + private String extractError(String body) { + if (StringUtils.isBlank(body)) { + return "无响应内容"; + } + try { + JSONObject error = JSON.parseObject(body); + return StringUtils.defaultIfBlank(error.getString("detail"), body); + } catch (Exception ignored) { + return StringUtils.abbreviate(body, 500); + } + } +} diff --git a/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisProperties.java b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisProperties.java new file mode 100644 index 0000000..30db743 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-llm/src/main/java/com/cmvr/llm/analysis/MediaAnalysisProperties.java @@ -0,0 +1,20 @@ +package com.cmvr.llm.analysis; + +import lombok.Data; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.stereotype.Component; + +@Data +@Component +@ConfigurationProperties(prefix = "media-analysis") +public class MediaAnalysisProperties { + + private String baseUrl = "http://192.168.28.10:14080"; + + private String apiKey; + + private int connectTimeoutMs = 5000; + + private int readTimeoutMs = 600000; +} + diff --git a/cmvr-iot-api/cmvr-iot-llm/src/test/java/com/cmvr/llm/analysis/MediaAnalysisClientImplTest.java b/cmvr-iot-api/cmvr-iot-llm/src/test/java/com/cmvr/llm/analysis/MediaAnalysisClientImplTest.java new file mode 100644 index 0000000..5faa710 --- /dev/null +++ b/cmvr-iot-api/cmvr-iot-llm/src/test/java/com/cmvr/llm/analysis/MediaAnalysisClientImplTest.java @@ -0,0 +1,17 @@ +package com.cmvr.llm.analysis; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class MediaAnalysisClientImplTest { + + @Test + void normalizesRawAndBearerPrefixedApiKeys() { + assertEquals("secret", MediaAnalysisClientImpl.normalizeApiKey("secret")); + assertEquals("secret", MediaAnalysisClientImpl.normalizeApiKey(" Bearer secret ")); + assertEquals("secret", MediaAnalysisClientImpl.normalizeApiKey("bearer secret")); + assertEquals("secret", MediaAnalysisClientImpl.normalizeApiKey("Bearer secret")); + assertEquals("", MediaAnalysisClientImpl.normalizeApiKey("Bearer ")); + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java index 5d36dc0..5b72d09 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/enums/ActionEnum.java @@ -84,6 +84,8 @@ public enum ActionEnum { INSPECTION_METER_RECOGNIZE("LLM", "INSPECTION_METER_RECOGNIZE", "巡检仪表读数识别"), AI_TTS("LLM", "AI_TTS", "tts语音播放"), GET_CURRENT_PAGE("LLM", "GET_CURRENT_PAGE", "获取当前页面名称"), + AUDIO_EVENT_CLASSIFY("LLM", "AUDIO_EVENT_CLASSIFY", "声音事件识别"), + VIDEO_ANALYZE("LLM", "VIDEO_ANALYZE", "视频智能分析"), // 触控交互 TI_PATH_SEARCH("EDGE", "TI_PATH_SEARCH", "路径搜索"), diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java new file mode 100644 index 0000000..551fbaf --- /dev/null +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateService.java @@ -0,0 +1,83 @@ +package com.cmvr.test.flow.runtime.operator.llm; + +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.common.exception.GlobalException; +import com.cmvr.llm.analysis.MediaAnalysisClient; +import com.cmvr.test.enums.ActionEnum; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; +import com.cmvr.test.flow.runtime.operator.FlowMediaParamResolver; +import lombok.RequiredArgsConstructor; +import org.apache.commons.lang3.StringUtils; +import org.springframework.stereotype.Service; + +import java.util.UUID; + +@Service +@RequiredArgsConstructor +public class LLMMediaAnalysisOperateService implements LLMOperateService { + + private final MediaAnalysisClient mediaAnalysisClient; + + @Override + public boolean supports(ActionEnum action) { + return action == ActionEnum.AUDIO_EVENT_CLASSIFY || action == ActionEnum.VIDEO_ANALYZE; + } + + @Override + public TaskNodeExecuteResult execute(TaskNodeExecuteMessage message) { + JSONObject input = message.getInputParams(); + boolean audio = message.getAction() == ActionEnum.AUDIO_EVENT_CLASSIFY; + String mediaParam = audio ? "audioUrl" : "videoUrl"; + String mediaUrl = FlowMediaParamResolver.lastUrl(input, mediaParam); + String profileCode = StringUtils.trimToNull(input.getString("profileCode")); + if (StringUtils.isBlank(mediaUrl)) { + throw new GlobalException("{}不能为空", audio ? "音频地址" : "视频地址"); + } + if (profileCode == null) { + throw new GlobalException("分析场景不能为空"); + } + + JSONObject context = new JSONObject(); + context.put("taskId", message.getTaskId()); + context.put("itemId", message.getItemId()); + context.put("flowInstId", message.getInstId()); + context.put("nodeId", message.getNodeId()); + context.put("trial", message.isTrial()); + + JSONObject options = new JSONObject(); + if (!audio) { + options.put("instruction", StringUtils.defaultString(input.getString("instruction"))); + options.put("analysisMode", StringUtils.defaultIfBlank(input.getString("analysisMode"), "AUTO")); + } + + JSONObject request = new JSONObject(); + request.put("requestId", buildRequestId(message)); + request.put("analysisType", audio ? "AUDIO_CLASSIFICATION" : "VIDEO_ANALYSIS"); + request.put("profileCode", profileCode); + request.put("mediaUrl", mediaUrl); + request.put("options", options); + request.put("context", context); + + JSONObject response = mediaAnalysisClient.analyze(request); + JSONObject output = response.getJSONObject("result"); + if (output == null) { + output = new JSONObject(); + } else { + output = new JSONObject(output); + } + output.put("analysisStatus", response.getString("status")); + output.put("analysisType", response.getString("analysisType")); + output.put("profileCode", response.getString("profileCode")); + output.put("analysisModel", response.getJSONObject("model")); + output.put("analysisTimingMs", response.getLong("timingMs")); + output.put("mediaUrl", mediaUrl); + return TaskNodeExecuteResult.success(output); + } + + private String buildRequestId(TaskNodeExecuteMessage message) { + String instance = StringUtils.defaultIfBlank(message.getInstId(), "single"); + String node = StringUtils.defaultIfBlank(message.getNodeId(), "node"); + return instance + "-" + node + "-" + UUID.randomUUID(); + } +} diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java index b846cc2..803f5b7 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowActionExecutorService.java @@ -311,6 +311,14 @@ public class FlowActionExecutorService { private String executeLLMAction(ActionEnum action, FlowActionRequestVO req) { log.info("LLM 执行动作: {}", action); + if (action == ActionEnum.AUDIO_EVENT_CLASSIFY || action == ActionEnum.VIDEO_ANALYZE) { + for (LLMOperateService service : llmOperateServices) { + if (service.supports(action)) { + return resultToString(service.execute(buildSingleNodeMessage(action, req.getPayload()))); + } + } + } + switch (action) { case TOUCH_COORDINATES: return "LLM 触控坐标 OK"; diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java new file mode 100644 index 0000000..71b2165 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/operator/llm/LLMMediaAnalysisOperateServiceTest.java @@ -0,0 +1,87 @@ +package com.cmvr.test.flow.runtime.operator.llm; + +import com.alibaba.fastjson2.JSONArray; +import com.alibaba.fastjson2.JSONObject; +import com.cmvr.llm.analysis.MediaAnalysisClient; +import com.cmvr.test.enums.ActionEnum; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage; +import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; +import org.junit.Test; + +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class LLMMediaAnalysisOperateServiceTest { + + @Test + public void mapsAudioArrayInputAndReturnsAnalysisMetadata() { + AtomicReference captured = new AtomicReference<>(); + MediaAnalysisClient client = request -> { + captured.set(request); + return new JSONObject() + .fluentPut("status", "SUCCEEDED") + .fluentPut("analysisType", "AUDIO_CLASSIFICATION") + .fluentPut("profileCode", "aima.power_state.v1") + .fluentPut("timingMs", 25L) + .fluentPut("model", new JSONObject().fluentPut("name", "audio-fingerprint")) + .fluentPut("result", new JSONObject() + .fluentPut("label", "POWER_ON") + .fluentPut("matched", true)); + }; + LLMMediaAnalysisOperateService service = new LLMMediaAnalysisOperateService(client); + JSONObject input = new JSONObject() + .fluentPut("profileCode", "aima.power_state.v1") + .fluentPut("audioUrl", new JSONArray() + .fluentAdd("https://files.example/first.wav") + .fluentAdd("https://files.example/latest.wav")); + + TaskNodeExecuteResult result = service.execute(message(ActionEnum.AUDIO_EVENT_CLASSIFY, input)); + + assertTrue(result.isSuccess()); + assertEquals("POWER_ON", result.getOutputParams().getString("label")); + assertEquals("https://files.example/latest.wav", captured.get().getString("mediaUrl")); + assertEquals("AUDIO_CLASSIFICATION", captured.get().getString("analysisType")); + assertEquals("SUCCEEDED", result.getOutputParams().getString("analysisStatus")); + } + + @Test + public void forwardsVideoInstruction() { + AtomicReference captured = new AtomicReference<>(); + MediaAnalysisClient client = request -> { + captured.set(request); + return new JSONObject() + .fluentPut("status", "SUCCEEDED") + .fluentPut("analysisType", "VIDEO_ANALYSIS") + .fluentPut("profileCode", "aima.power_video.v1") + .fluentPut("timingMs", 30L) + .fluentPut("model", new JSONObject().fluentPut("name", "qwen3-vl:32b")) + .fluentPut("result", new JSONObject().fluentPut("summary", "ok")); + }; + LLMMediaAnalysisOperateService service = new LLMMediaAnalysisOperateService(client); + JSONObject input = new JSONObject() + .fluentPut("profileCode", "aima.power_video.v1") + .fluentPut("videoUrl", "https://files.example/video.mp4") + .fluentPut("analysisMode", "FAST") + .fluentPut("instruction", "检查仪表是否点亮"); + + service.execute(message(ActionEnum.VIDEO_ANALYZE, input)); + + assertEquals("VIDEO_ANALYSIS", captured.get().getString("analysisType")); + assertEquals("检查仪表是否点亮", + captured.get().getJSONObject("options").getString("instruction")); + assertEquals("FAST", captured.get().getJSONObject("options").getString("analysisMode")); + } + + private TaskNodeExecuteMessage message(ActionEnum action, JSONObject input) { + TaskNodeExecuteMessage message = new TaskNodeExecuteMessage(); + message.setAction(action); + message.setInstId("inst-1"); + message.setTaskId("task-1"); + message.setItemId("item-1"); + message.setNodeId("node-1"); + message.setInputParams(input); + return message; + } +}