Compare commits
2 Commits
05992d1980
...
c4a9f29780
| Author | SHA1 | Date | |
|---|---|---|---|
| c4a9f29780 | |||
| 6265db01c7 |
@ -132,3 +132,9 @@ flowise:
|
|||||||
abort: http://192.168.0.108:3000/api/v1/chatmessage/abort/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/
|
query: http://192.168.0.108:3000/api/v1/executions/
|
||||||
api-key: pg61JW6W_GXyqmqSoURJ4mlSRrCMRzGLBJli_w3BFkg
|
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
|
||||||
|
|||||||
@ -132,3 +132,9 @@ flowise:
|
|||||||
abort: http://192.168.0.108:3000/api/v1/chatmessage/abort/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/
|
query: http://192.168.0.108:3000/api/v1/executions/
|
||||||
api-key: pg61JW6W_GXyqmqSoURJ4mlSRrCMRzGLBJli_w3BFkg
|
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
|
||||||
|
|||||||
@ -132,3 +132,9 @@ flowise:
|
|||||||
abort: http://192.168.0.108:3000/api/v1/chatmessage/abort/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/
|
query: http://192.168.0.108:3000/api/v1/executions/
|
||||||
api-key: pg61JW6W_GXyqmqSoURJ4mlSRrCMRzGLBJli_w3BFkg
|
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
|
||||||
|
|||||||
@ -342,6 +342,13 @@ public class EdgeSystemServiceImpl implements EdgeSystemService {
|
|||||||
.setHeader(EdgeCommonUtil.buildRequest(""))
|
.setHeader(EdgeCommonUtil.buildRequest(""))
|
||||||
.build();
|
.build();
|
||||||
SystemCommand.StopAllCommand.Feedback feedback = stub.stopAll(build);
|
SystemCommand.StopAllCommand.Feedback feedback = stub.stopAll(build);
|
||||||
|
if (feedback == null || !feedback.hasHeader()) {
|
||||||
|
throw new GlobalException("停止机器人全部设备失败: 响应为空");
|
||||||
|
}
|
||||||
|
if (!feedback.getHeader().getSuccess()) {
|
||||||
|
throw new GlobalException("停止机器人全部设备失败: "
|
||||||
|
+ feedback.getHeader().getErrorMessage());
|
||||||
|
}
|
||||||
return JSON.toJSONString(feedback.getHeader());
|
return JSON.toJSONString(feedback.getHeader());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -30,6 +30,12 @@
|
|||||||
<version>1.0.49</version>
|
<version>1.0.49</version>
|
||||||
</dependency>
|
</dependency>
|
||||||
|
|
||||||
|
<dependency>
|
||||||
|
<groupId>org.springframework.boot</groupId>
|
||||||
|
<artifactId>spring-boot-starter-test</artifactId>
|
||||||
|
<scope>test</scope>
|
||||||
|
</dependency>
|
||||||
|
|
||||||
</dependencies>
|
</dependencies>
|
||||||
|
|
||||||
</project>
|
</project>
|
||||||
@ -0,0 +1,9 @@
|
|||||||
|
package com.cmvr.llm.analysis;
|
||||||
|
|
||||||
|
import com.alibaba.fastjson2.JSONObject;
|
||||||
|
|
||||||
|
public interface MediaAnalysisClient {
|
||||||
|
|
||||||
|
JSONObject analyze(JSONObject request);
|
||||||
|
}
|
||||||
|
|
||||||
@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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;
|
||||||
|
}
|
||||||
|
|
||||||
@ -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 "));
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -84,6 +84,8 @@ public enum ActionEnum {
|
|||||||
INSPECTION_METER_RECOGNIZE("LLM", "INSPECTION_METER_RECOGNIZE", "巡检仪表读数识别"),
|
INSPECTION_METER_RECOGNIZE("LLM", "INSPECTION_METER_RECOGNIZE", "巡检仪表读数识别"),
|
||||||
AI_TTS("LLM", "AI_TTS", "tts语音播放"),
|
AI_TTS("LLM", "AI_TTS", "tts语音播放"),
|
||||||
GET_CURRENT_PAGE("LLM", "GET_CURRENT_PAGE", "获取当前页面名称"),
|
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", "路径搜索"),
|
TI_PATH_SEARCH("EDGE", "TI_PATH_SEARCH", "路径搜索"),
|
||||||
|
|||||||
@ -84,6 +84,9 @@ public class TaskContext {
|
|||||||
*/
|
*/
|
||||||
private volatile boolean paused = false;
|
private volatile boolean paused = false;
|
||||||
|
|
||||||
|
/** 暂停时的设备停止和旧节点退出是否已经完成。 */
|
||||||
|
private volatile boolean pauseSettled = true;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 是否已被强行终止
|
* 是否已被强行终止
|
||||||
*/
|
*/
|
||||||
|
|||||||
@ -96,10 +96,14 @@ public class TaskThreadRegistry {
|
|||||||
/**
|
/**
|
||||||
* 移除节点上下文(节点执行成功时调用)
|
* 移除节点上下文(节点执行成功时调用)
|
||||||
*/
|
*/
|
||||||
public void unregisterNodeContext(String instId, String itemId, String nodeId) {
|
public void unregisterNodeContext(String instId, String itemId, String nodeId,
|
||||||
|
TaskNodeExecuteContext expectedContext) {
|
||||||
String key = TaskKeyBuilder.buildKey(instId, itemId, nodeId);
|
String key = TaskKeyBuilder.buildKey(instId, itemId, nodeId);
|
||||||
nodeContextMap.remove(key);
|
if (nodeContextMap.remove(key, expectedContext)) {
|
||||||
log.info("移除节点上下文: key={}", key);
|
log.info("移除节点上下文: key={}", key);
|
||||||
|
} else {
|
||||||
|
log.debug("忽略过期节点上下文注销: key={}", key);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@ -24,6 +24,7 @@ import java.util.ArrayList;
|
|||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Collections;
|
import java.util.Collections;
|
||||||
import java.util.concurrent.CountDownLatch;
|
import java.util.concurrent.CountDownLatch;
|
||||||
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@ -35,6 +36,8 @@ import java.util.stream.Collectors;
|
|||||||
@RequiredArgsConstructor
|
@RequiredArgsConstructor
|
||||||
public class FlowControlService {
|
public class FlowControlService {
|
||||||
|
|
||||||
|
private static final long PAUSE_SETTLE_TIMEOUT_SECONDS = 15L;
|
||||||
|
|
||||||
@Resource(name = "threadPoolTaskExecutor")
|
@Resource(name = "threadPoolTaskExecutor")
|
||||||
private final ThreadPoolTaskExecutor executor;
|
private final ThreadPoolTaskExecutor executor;
|
||||||
private final TaskInstHolder taskInstHolder;
|
private final TaskInstHolder taskInstHolder;
|
||||||
@ -80,15 +83,22 @@ public class FlowControlService {
|
|||||||
if (ctx.isPaused()) {
|
if (ctx.isPaused()) {
|
||||||
throw new GlobalException("任务已处于暂停状态");
|
throw new GlobalException("任务已处于暂停状态");
|
||||||
}
|
}
|
||||||
// 异步停止终端
|
|
||||||
stopAllRobots(ctx);
|
|
||||||
ctx.setPaused(true);
|
ctx.setPaused(true);
|
||||||
|
ctx.setPauseSettled(false);
|
||||||
ctx.setTerminalStatus(TerminalStatusEnum.RUNNING);
|
ctx.setTerminalStatus(TerminalStatusEnum.RUNNING);
|
||||||
ctx.setStatus(TaskStatusEnum.PAUSED);
|
ctx.setStatus(TaskStatusEnum.PAUSED);
|
||||||
ctx.updateTimestamp();
|
ctx.updateTimestamp();
|
||||||
|
List<TaskNodeExecuteContext> interruptedNodes = taskThreadRegistry.getPendingNodes(instId);
|
||||||
|
taskThreadRegistry.interruptAll(instId);
|
||||||
|
// StopAll 是阻塞式调用。只有设备侧停止完成,暂停接口才返回,避免恢复请求和旧动作并发。
|
||||||
|
boolean robotsStopped = stopAllRobotsAndWait(ctx);
|
||||||
|
boolean nodesFinished = waitForInterruptedNodes(instId, interruptedNodes);
|
||||||
|
ctx.setPauseSettled(robotsStopped && nodesFinished);
|
||||||
|
if (!ctx.isPauseSettled()) {
|
||||||
|
log.warn("任务已进入暂停状态,但设备停止仍在收敛: instId={}", instId);
|
||||||
|
}
|
||||||
// 统一记录日志 + 设置上下文状态 + 数据库状态
|
// 统一记录日志 + 设置上下文状态 + 数据库状态
|
||||||
taskInstHolder.syncStatus(instId, TaskStatusEnum.PAUSED);
|
taskInstHolder.syncStatus(instId, TaskStatusEnum.PAUSED);
|
||||||
taskThreadRegistry.interruptAll(instId);
|
|
||||||
// 被中断的节点 状态修改成 PAUSED
|
// 被中断的节点 状态修改成 PAUSED
|
||||||
for (TaskNodeExecuteContext pendingNode : taskThreadRegistry.getPendingNodes(instId)) {
|
for (TaskNodeExecuteContext pendingNode : taskThreadRegistry.getPendingNodes(instId)) {
|
||||||
String nodeId = pendingNode.getNode().getNodeId();
|
String nodeId = pendingNode.getNode().getNodeId();
|
||||||
@ -107,6 +117,15 @@ public class FlowControlService {
|
|||||||
if (!ctx.isPaused()) {
|
if (!ctx.isPaused()) {
|
||||||
throw new GlobalException("当前任务未处于暂停状态,无法继续运行");
|
throw new GlobalException("当前任务未处于暂停状态,无法继续运行");
|
||||||
}
|
}
|
||||||
|
if (!ctx.isPauseSettled()) {
|
||||||
|
List<TaskNodeExecuteContext> interruptedNodes = taskThreadRegistry.getPendingNodes(instId);
|
||||||
|
boolean robotsStopped = stopAllRobotsAndWait(ctx);
|
||||||
|
boolean nodesFinished = waitForInterruptedNodes(instId, interruptedNodes);
|
||||||
|
ctx.setPauseSettled(robotsStopped && nodesFinished);
|
||||||
|
if (!ctx.isPauseSettled()) {
|
||||||
|
throw new GlobalException("设备停止尚未完成,请稍后重新恢复");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
ctx.setPaused(false);
|
ctx.setPaused(false);
|
||||||
|
|
||||||
@ -207,4 +226,48 @@ public class FlowControlService {
|
|||||||
.distinct()
|
.distinct()
|
||||||
.forEach(robotId -> executor.execute(() -> edgeSystemService.stopAll(robotId)));
|
.forEach(robotId -> executor.execute(() -> edgeSystemService.stopAll(robotId)));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private boolean stopAllRobotsAndWait(TaskContext context) {
|
||||||
|
List<String> robotIds = context.getRobotIds() == null || context.getRobotIds().isEmpty()
|
||||||
|
? Collections.singletonList(context.getRobotId())
|
||||||
|
: context.getRobotIds();
|
||||||
|
boolean success = true;
|
||||||
|
for (String robotId : robotIds.stream()
|
||||||
|
.filter(robotId -> robotId != null && !robotId.isBlank())
|
||||||
|
.distinct()
|
||||||
|
.collect(Collectors.toList())) {
|
||||||
|
try {
|
||||||
|
edgeSystemService.stopAll(robotId);
|
||||||
|
} catch (RuntimeException exception) {
|
||||||
|
success = false;
|
||||||
|
log.warn("暂停机器人设备失败,恢复时将重试: robotId={}, error={}",
|
||||||
|
robotId, exception.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return success;
|
||||||
|
}
|
||||||
|
|
||||||
|
private boolean waitForInterruptedNodes(String instId, List<TaskNodeExecuteContext> interruptedNodes) {
|
||||||
|
long deadlineNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(PAUSE_SETTLE_TIMEOUT_SECONDS);
|
||||||
|
for (TaskNodeExecuteContext pendingNode : interruptedNodes) {
|
||||||
|
if (pendingNode.getNode().getNodeType() == NodeTypeEnum.LOOP) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
long remainingNanos = deadlineNanos - System.nanoTime();
|
||||||
|
if (remainingNanos <= 0) {
|
||||||
|
log.warn("暂停等待节点退出超时: instId={}, nodeId={}", instId, pendingNode.getNode().getNodeId());
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
if (!pendingNode.awaitExecutionFinished(remainingNanos, TimeUnit.NANOSECONDS)) {
|
||||||
|
log.warn("暂停等待节点退出超时: instId={}, nodeId={}", instId, pendingNode.getNode().getNodeId());
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
} catch (InterruptedException exception) {
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
|
throw new GlobalException("暂停任务时等待设备动作结束被中断");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -118,7 +118,9 @@ public class FlowItemExecutor {
|
|||||||
log.info("当前执行的是 {} 节点,迭代路径={}", node.getNodeName(), iterations);
|
log.info("当前执行的是 {} 节点,迭代路径={}", node.getNodeName(), iterations);
|
||||||
|
|
||||||
//注册当前线程
|
//注册当前线程
|
||||||
registerThread(graph, node, startNodeId, startInputDefs, rootMessage, onFinished, iterations, instId, itemId, nodeId);
|
TaskNodeExecuteContext nodeExecutionContext = registerThread(
|
||||||
|
graph, node, startNodeId, startInputDefs, rootMessage, onFinished,
|
||||||
|
iterations, instId, itemId, nodeId);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
// 创建带上下文的执行消息
|
// 创建带上下文的执行消息
|
||||||
@ -172,10 +174,11 @@ public class FlowItemExecutor {
|
|||||||
}
|
}
|
||||||
throw e;
|
throw e;
|
||||||
} finally {
|
} finally {
|
||||||
|
nodeExecutionContext.markExecutionFinished();
|
||||||
// 如果不是暂停全中断
|
// 如果不是暂停全中断
|
||||||
TaskContext ctx = taskInstHolder.getContext(instId);
|
TaskContext ctx = taskInstHolder.getContext(instId);
|
||||||
if (ctx == null || !ctx.isPaused()) {
|
if (ctx == null || !ctx.isPaused()) {
|
||||||
taskThreadRegistry.unregisterNodeContext(instId, itemId, nodeId);
|
taskThreadRegistry.unregisterNodeContext(instId, itemId, nodeId, nodeExecutionContext);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -286,7 +289,7 @@ public class FlowItemExecutor {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
private void registerThread(FlowGraph graph, FlowNodeWrapper node, String startNodeId, List<FlowParamDef> startInputDefs,
|
private TaskNodeExecuteContext registerThread(FlowGraph graph, FlowNodeWrapper node, String startNodeId, List<FlowParamDef> startInputDefs,
|
||||||
TaskNodeExecuteMessage rootMessage, Runnable onFinished, List<Integer> iterations, String instId, String itemId, String nodeId) {
|
TaskNodeExecuteMessage rootMessage, Runnable onFinished, List<Integer> iterations, String instId, String itemId, String nodeId) {
|
||||||
TaskNodeExecuteContext nodeExecuteContext = new TaskNodeExecuteContext();
|
TaskNodeExecuteContext nodeExecuteContext = new TaskNodeExecuteContext();
|
||||||
nodeExecuteContext.setThread(Thread.currentThread());
|
nodeExecuteContext.setThread(Thread.currentThread());
|
||||||
@ -299,6 +302,7 @@ public class FlowItemExecutor {
|
|||||||
// nodeExecuteContext.setLoopIteration(loopIteration);
|
// nodeExecuteContext.setLoopIteration(loopIteration);
|
||||||
nodeExecuteContext.setIterations(new ArrayList<>(iterations));
|
nodeExecuteContext.setIterations(new ArrayList<>(iterations));
|
||||||
taskThreadRegistry.registerNodeContext(instId, itemId, nodeId, nodeExecuteContext);
|
taskThreadRegistry.registerNodeContext(instId, itemId, nodeId, nodeExecuteContext);
|
||||||
|
return nodeExecuteContext;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@ -7,12 +7,16 @@ import lombok.Data;
|
|||||||
|
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
import java.util.concurrent.CountDownLatch;
|
||||||
|
import java.util.concurrent.TimeUnit;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 节点执行上下文,用于保存线程与流程信息
|
* 节点执行上下文,用于保存线程与流程信息
|
||||||
*/
|
*/
|
||||||
@Data
|
@Data
|
||||||
public class TaskNodeExecuteContext {
|
public class TaskNodeExecuteContext {
|
||||||
|
private final CountDownLatch executionFinished = new CountDownLatch(1);
|
||||||
|
|
||||||
private Thread thread;
|
private Thread thread;
|
||||||
|
|
||||||
private FlowGraph graph;
|
private FlowGraph graph;
|
||||||
@ -24,4 +28,12 @@ public class TaskNodeExecuteContext {
|
|||||||
private int loopIteration;
|
private int loopIteration;
|
||||||
private Object loopArray;
|
private Object loopArray;
|
||||||
private List<Integer> iterations = new ArrayList<>();
|
private List<Integer> iterations = new ArrayList<>();
|
||||||
|
|
||||||
|
public void markExecutionFinished() {
|
||||||
|
executionFinished.countDown();
|
||||||
|
}
|
||||||
|
|
||||||
|
public boolean awaitExecutionFinished(long timeout, TimeUnit unit) throws InterruptedException {
|
||||||
|
return executionFinished.await(timeout, unit);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -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();
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -311,6 +311,14 @@ public class FlowActionExecutorService {
|
|||||||
private String executeLLMAction(ActionEnum action, FlowActionRequestVO req) {
|
private String executeLLMAction(ActionEnum action, FlowActionRequestVO req) {
|
||||||
log.info("LLM 执行动作: {}", action);
|
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) {
|
switch (action) {
|
||||||
case TOUCH_COORDINATES:
|
case TOUCH_COORDINATES:
|
||||||
return "LLM 触控坐标 OK";
|
return "LLM 触控坐标 OK";
|
||||||
|
|||||||
@ -0,0 +1,24 @@
|
|||||||
|
package com.cmvr.test.flow.context;
|
||||||
|
|
||||||
|
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteContext;
|
||||||
|
import org.junit.Test;
|
||||||
|
|
||||||
|
import static org.junit.Assert.assertSame;
|
||||||
|
|
||||||
|
public class TaskThreadRegistryTest {
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void staleExecutionCannotUnregisterResumedNodeContext() {
|
||||||
|
TaskThreadRegistry registry = new TaskThreadRegistry();
|
||||||
|
TaskNodeExecuteContext interrupted = new TaskNodeExecuteContext();
|
||||||
|
TaskNodeExecuteContext resumed = new TaskNodeExecuteContext();
|
||||||
|
interrupted.setThread(Thread.currentThread());
|
||||||
|
resumed.setThread(Thread.currentThread());
|
||||||
|
|
||||||
|
registry.registerNodeContext("inst-1", "item-1", "node-1", interrupted);
|
||||||
|
registry.registerNodeContext("inst-1", "item-1", "node-1", resumed);
|
||||||
|
registry.unregisterNodeContext("inst-1", "item-1", "node-1", interrupted);
|
||||||
|
|
||||||
|
assertSame(resumed, registry.getNodeContext("inst-1", "item-1", "node-1"));
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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<JSONObject> 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<JSONObject> 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;
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Reference in New Issue
Block a user