feat(llm): 添加通用图片分析功能并优化边缘系统安全状态恢复

- 新增 IMAGE_ANALYZE 动作类型支持通用图片分析
- 实现图片分析操作服务支持单张及多图片分析
- 添加图片分析的提示词、参考图片和调优参数配置
- 集成 MinIO 服务上传功能支持唯一文件名生成
- 添加 FlowMediaParamResolver 支持 URL 数组解析
- 实现边缘系统恢复运行状态的新 API 接口
- 标记旧的安全状态恢复接口为已弃用
- 更新媒体分析操作服务支持详细的输出选项
- 添加图片分析相关的单元测试验证功能
This commit is contained in:
lixiaolong 2026-08-19 15:00:49 +08:00
parent 937426fcfc
commit 53ed826480
23 changed files with 3250 additions and 93 deletions

View File

@ -63,7 +63,18 @@ public class EdgeSystemController extends BaseController {
return success(JSON.toJSONString(edgeSystemService.getSafetyState(robotId)));
}
@ApiOperation("Clear recoverable edge safety errors")
@ApiOperation("Restore edge operational state")
@PostMapping("/restoreOperationalState")
public AjaxResult restoreOperationalState(@RequestBody @Validated EdgeSafetyRecoveryVO request) {
return success(JSON.toJSONString(edgeSystemService.restoreOperationalState(
request.getRobotId(), request.getReason(), request.getTimeoutMs())));
}
/**
* @deprecated Use /restoreOperationalState.
*/
@Deprecated
@ApiOperation("[Deprecated] Clear recoverable edge safety errors")
@PostMapping("/recoverSafetyState")
public AjaxResult recoverSafetyState(@RequestBody @Validated EdgeSafetyRecoveryVO request) {
return success(JSON.toJSONString(edgeSystemService.recoverSafetyState(

View File

@ -31,6 +31,7 @@ public final class AimaAnalysisVerdict {
}
private static boolean isAnalysisAction(ActionEnum action) {
return action == ActionEnum.AUDIO_EVENT_CLASSIFY || action == ActionEnum.VIDEO_ANALYZE;
return action == ActionEnum.AUDIO_EVENT_CLASSIFY
|| action == ActionEnum.VIDEO_ANALYZE;
}
}

View File

@ -20,6 +20,8 @@ public class AimaAnalysisVerdictTest {
assertEquals(AimaAnalysisVerdict.NOT_PASSED,
AimaAnalysisVerdict.resolve(ActionEnum.VIDEO_ANALYZE,
new JSONObject().fluentPut("analysisStatus", "FAILED")));
assertNull(AimaAnalysisVerdict.resolve(ActionEnum.IMAGE_ANALYZE,
new JSONObject().fluentPut("result", "12.6")));
assertNull(AimaAnalysisVerdict.resolve(ActionEnum.SLEEP,
new JSONObject().fluentPut("passed", true)));
}

View File

@ -30,6 +30,13 @@ public interface EdgeSystemService {
SafetyCommand.GetSafetyStateCommand.Feedback getSafetyState(String robotId);
SafetyCommand.RestoreOperationalStateCommand.Feedback restoreOperationalState(
String robotId, String reason, Integer timeoutMs);
/**
* @deprecated Use {@link #restoreOperationalState(String, String, Integer)}.
*/
@Deprecated
SafetyCommand.RecoverSafetyStateCommand.Feedback recoverSafetyState(
String robotId, String reason, Integer timeoutMs);

View File

@ -135,6 +135,39 @@ public class EdgeSystemServiceImpl implements EdgeSystemService {
}
@Override
public SafetyCommand.RestoreOperationalStateCommand.Feedback restoreOperationalState(
String robotId, String reason, Integer timeoutMs) {
if (robotId == null || robotId.trim().isEmpty()) {
throw new GlobalException("机器人ID不能为空");
}
SafetyCommand.RestoreOperationalStateCommand.Request request =
buildRestoreOperationalStateRequest(reason, timeoutMs);
SystemServiceGrpc.SystemServiceBlockingStub stub = grpcServiceManager.getGrpcClient(
robotId.trim(), SystemServiceGrpc.SystemServiceBlockingStub.class);
SafetyCommand.RestoreOperationalStateCommand.Feedback feedback = stub
.withDeadlineAfter(request.getTimeoutMs() + 5000L, TimeUnit.MILLISECONDS)
.restoreOperationalState(request);
requireSuccessfulHeader(feedback == null || !feedback.hasHeader() ? null : feedback.getHeader(),
"恢复机器人运行状态");
if (!isSuccessfulRecoveryResult(feedback.getResult())) {
throw new GlobalException("恢复机器人运行状态未完成: " + feedback.getResult().name());
}
return feedback;
}
private SafetyCommand.RestoreOperationalStateCommand.Request buildRestoreOperationalStateRequest(
String reason, Integer timeoutMs) {
int effectiveTimeoutMs = timeoutMs == null ? 15000 : Math.max(1000, Math.min(timeoutMs, 120000));
return SafetyCommand.RestoreOperationalStateCommand.Request.newBuilder()
.setRecoveryId(UUID.randomUUID().toString())
.setReason(reason == null || reason.trim().isEmpty()
? "Platform manual operational state restore" : reason.trim())
.setTimeoutMs(effectiveTimeoutMs)
.build();
}
@Override
@Deprecated
public SafetyCommand.RecoverSafetyStateCommand.Feedback recoverSafetyState(
String robotId, String reason, Integer timeoutMs) {
SafetyCommand.GetSafetyStateCommand.Feedback current = getSafetyState(robotId);

View File

@ -263,6 +263,37 @@ public final class SystemServiceGrpc {
return getRecoverSafetyStateMethod;
}
private static volatile io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request,
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback> getRestoreOperationalStateMethod;
@io.grpc.stub.annotations.RpcMethod(
fullMethodName = SERVICE_NAME + '/' + "RestoreOperationalState",
requestType = cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request.class,
responseType = cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback.class,
methodType = io.grpc.MethodDescriptor.MethodType.UNARY)
public static io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request,
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback> getRestoreOperationalStateMethod() {
io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request, cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback> getRestoreOperationalStateMethod;
if ((getRestoreOperationalStateMethod = SystemServiceGrpc.getRestoreOperationalStateMethod) == null) {
synchronized (SystemServiceGrpc.class) {
if ((getRestoreOperationalStateMethod = SystemServiceGrpc.getRestoreOperationalStateMethod) == null) {
SystemServiceGrpc.getRestoreOperationalStateMethod = getRestoreOperationalStateMethod =
io.grpc.MethodDescriptor.<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request, cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback>newBuilder()
.setType(io.grpc.MethodDescriptor.MethodType.UNARY)
.setFullMethodName(generateFullMethodName(SERVICE_NAME, "RestoreOperationalState"))
.setSampledToLocalTracing(true)
.setRequestMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request.getDefaultInstance()))
.setResponseMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback.getDefaultInstance()))
.setSchemaDescriptor(new SystemServiceMethodDescriptorSupplier("RestoreOperationalState"))
.build();
}
}
}
return getRestoreOperationalStateMethod;
}
/**
* Creates a new async stub that supports all call types for the service
*/
@ -367,6 +398,13 @@ public final class SystemServiceGrpc {
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getRecoverSafetyStateMethod(), responseObserver);
}
/**
*/
public void restoreOperationalState(cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request request,
io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback> responseObserver) {
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getRestoreOperationalStateMethod(), responseObserver);
}
@java.lang.Override public final io.grpc.ServerServiceDefinition bindService() {
return io.grpc.ServerServiceDefinition.builder(getServiceDescriptor())
.addMethod(
@ -425,6 +463,13 @@ public final class SystemServiceGrpc {
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request,
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback>(
this, METHODID_RECOVER_SAFETY_STATE)))
.addMethod(
getRestoreOperationalStateMethod(),
io.grpc.stub.ServerCalls.asyncUnaryCall(
new MethodHandlers<
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request,
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback>(
this, METHODID_RESTORE_OPERATIONAL_STATE)))
.build();
}
}
@ -506,6 +551,14 @@ public final class SystemServiceGrpc {
io.grpc.stub.ClientCalls.asyncUnaryCall(
getChannel().newCall(getRecoverSafetyStateMethod(), getCallOptions()), request, responseObserver);
}
/**
*/
public void restoreOperationalState(cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request request,
io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback> responseObserver) {
io.grpc.stub.ClientCalls.asyncUnaryCall(
getChannel().newCall(getRestoreOperationalStateMethod(), getCallOptions()), request, responseObserver);
}
}
/**
@ -577,6 +630,13 @@ public final class SystemServiceGrpc {
return io.grpc.stub.ClientCalls.blockingUnaryCall(
getChannel(), getRecoverSafetyStateMethod(), getCallOptions(), request);
}
/**
*/
public cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback restoreOperationalState(cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request request) {
return io.grpc.stub.ClientCalls.blockingUnaryCall(
getChannel(), getRestoreOperationalStateMethod(), getCallOptions(), request);
}
}
/**
@ -656,6 +716,14 @@ public final class SystemServiceGrpc {
return io.grpc.stub.ClientCalls.futureUnaryCall(
getChannel().newCall(getRecoverSafetyStateMethod(), getCallOptions()), request);
}
/**
*/
public com.google.common.util.concurrent.ListenableFuture<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback> restoreOperationalState(
cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request request) {
return io.grpc.stub.ClientCalls.futureUnaryCall(
getChannel().newCall(getRestoreOperationalStateMethod(), getCallOptions()), request);
}
}
private static final int METHODID_GET_SYSTEM_INFO = 0;
@ -666,6 +734,7 @@ public final class SystemServiceGrpc {
private static final int METHODID_EXECUTE_ACTION_QUEUE = 5;
private static final int METHODID_GET_SAFETY_STATE = 6;
private static final int METHODID_RECOVER_SAFETY_STATE = 7;
private static final int METHODID_RESTORE_OPERATIONAL_STATE = 8;
private static final class MethodHandlers<Req, Resp> implements
io.grpc.stub.ServerCalls.UnaryMethod<Req, Resp>,
@ -716,6 +785,10 @@ public final class SystemServiceGrpc {
serviceImpl.recoverSafetyState((cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback>) responseObserver);
break;
case METHODID_RESTORE_OPERATIONAL_STATE:
serviceImpl.restoreOperationalState((cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RestoreOperationalStateCommand.Feedback>) responseObserver);
break;
default:
throw new AssertionError();
}
@ -785,6 +858,7 @@ public final class SystemServiceGrpc {
.addMethod(getExecuteActionQueueMethod())
.addMethod(getGetSafetyStateMethod())
.addMethod(getRecoverSafetyStateMethod())
.addMethod(getRestoreOperationalStateMethod())
.build();
}
}

View File

@ -25,7 +25,7 @@ public final class SystemServiceOuterClass {
java.lang.String[] descriptorData = {
"\n\035cmvr/api/system_service.proto\022\010cmvr.ap" +
"i\032\035cmvr/api/system_command.proto\032\035cmvr/a" +
"pi/safety_command.proto2\263\006\n\rSystemServic" +
"pi/safety_command.proto2\266\007\n\rSystemServic" +
"e\022b\n\rGetSystemInfo\022&.cmvr.api.GetSystemI" +
"nfoCommand.Request\032\'.cmvr.api.GetSystemI" +
"nfoCommand.Feedback\"\000\022h\n\017GetSystemStatus" +
@ -46,7 +46,10 @@ public final class SystemServiceOuterClass {
"Feedback\"\000\022q\n\022RecoverSafetyState\022+.cmvr." +
"api.RecoverSafetyStateCommand.Request\032,." +
"cmvr.api.RecoverSafetyStateCommand.Feedb" +
"ack\"\000b\006proto3"
"ack\"\000\022\200\001\n\027RestoreOperationalState\0220.cmvr" +
".api.RestoreOperationalStateCommand.Requ" +
"est\0321.cmvr.api.RestoreOperationalStateCo" +
"mmand.Feedback\"\000b\006proto3"
};
descriptor = com.google.protobuf.Descriptors.FileDescriptor
.internalBuildGeneratedFileFrom(descriptorData,

View File

@ -184,3 +184,21 @@ message RecoverSafetyStateCommand {
string recovery_id = 7;
}
}
message RestoreOperationalStateCommand {
message Request {
string recovery_id = 1;
string reason = 2;
uint32 timeout_ms = 3;
}
message Feedback {
CommandHeader.Feedback header = 1;
SafetyOperationResult result = 2;
uint64 previous_safety_epoch = 3;
uint64 current_safety_epoch = 4;
SystemAdmissionState system_state = 5;
repeated SafetyOperationTargetResult targets = 6;
string recovery_id = 7;
}
}

View File

@ -18,5 +18,8 @@ service SystemService {
rpc ExecuteActionQueue(ActionQueueCommand.Request) returns (ActionQueueCommand.Feedback) {}
rpc GetSafetyState(GetSafetyStateCommand.Request) returns (GetSafetyStateCommand.Feedback) {}
rpc RecoverSafetyState(RecoverSafetyStateCommand.Request) returns (RecoverSafetyStateCommand.Feedback) {}
rpc RecoverSafetyState(RecoverSafetyStateCommand.Request) returns (RecoverSafetyStateCommand.Feedback) {
option deprecated = true;
}
rpc RestoreOperationalState(RestoreOperationalStateCommand.Request) returns (RestoreOperationalStateCommand.Feedback) {}
}

View File

@ -156,6 +156,13 @@
<version>${swagger.annotations.version}</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.13.2</version>
<scope>test</scope>
</dependency>
</dependencies>
</project>

View File

@ -25,6 +25,7 @@ import java.util.Date;
import java.util.HashMap;
import java.util.Locale;
import java.util.Map;
import java.util.UUID;
@Slf4j
@Component
@ -96,16 +97,12 @@ public class MinioService {
if (StringUtils.isBlank(originalFilename)) {
throw new RuntimeException();
}
String fileName = System.currentTimeMillis() + "";
if (originalFilename.lastIndexOf(".") >= 0) {
fileName += originalFilename.substring(originalFilename.lastIndexOf("."));
}
String fileName = buildUniqueFileName(originalFilename);
String objectName = DateUtil.format(new Date(), DatePattern.PURE_DATE_PATTERN) + "/" + fileName;
try {
InputStream stream = new ByteArrayInputStream(file.getBytes());
PutObjectArgs objectArgs = PutObjectArgs.builder().bucket(bucketName).object(path + "/" + objectName)
.stream(stream, file.getSize(), -1).contentType(file.getContentType()).build();
//文件名称相同会覆盖
minioClient.putObject(objectArgs);
} catch (Exception e) {
log.error("上传图片失败,", e);
@ -121,11 +118,7 @@ public class MinioService {
throw new RuntimeException("文件名为空");
}
// 生成文件名,使用当前时间戳
String fileName = System.currentTimeMillis() + "";
if (originalFilename.lastIndexOf(".") >= 0) {
fileName += originalFilename.substring(originalFilename.lastIndexOf("."));
}
String fileName = buildUniqueFileName(originalFilename);
// 构建文件在 MinIO 中的路径
String objectName = DateUtil.format(new Date(), DatePattern.PURE_DATE_PATTERN) + "/" + fileName;
@ -178,13 +171,12 @@ public class MinioService {
*/
public String upload(String bucketName, String path, String name, byte[] fileByte, String contentType) {
String fileName = System.currentTimeMillis() + name.substring(name.lastIndexOf("."));
String fileName = buildUniqueFileName(name);
String objectName = DateUtil.format(new Date(), DatePattern.PURE_DATE_PATTERN) + "/" + fileName;
try {
InputStream stream = new ByteArrayInputStream(fileByte);
PutObjectArgs objectArgs = PutObjectArgs.builder().bucket(bucketName).object(path + "/" + objectName)
.stream(stream, fileByte.length, -1).contentType(contentType).build();
//文件名称相同会覆盖
minioClient.putObject(objectArgs);
} catch (Exception e) {
log.error("上传文件失败,", e);
@ -193,6 +185,30 @@ public class MinioService {
return objectName;
}
/**
* Generates a collision-resistant object file name while retaining the source extension.
*/
static String buildUniqueFileName(String originalFilename) {
String extension = "";
if (StringUtils.isNotBlank(originalFilename)) {
int separatorIndex = Math.max(originalFilename.lastIndexOf('/'), originalFilename.lastIndexOf('\\'));
int extensionIndex = originalFilename.lastIndexOf('.');
if (extensionIndex > separatorIndex) {
extension = originalFilename.substring(extensionIndex);
}
}
String uuid = UUID.randomUUID().toString().replace("-", "");
return System.currentTimeMillis() + "-" + uuid + extension;
}
static String buildUniqueObjectName(String objectName) {
String normalized = StringUtils.defaultString(objectName).replace('\\', '/');
int separatorIndex = normalized.lastIndexOf('/');
String directory = separatorIndex >= 0 ? normalized.substring(0, separatorIndex + 1) : "";
String originalFilename = separatorIndex >= 0 ? normalized.substring(separatorIndex + 1) : normalized;
return directory + buildUniqueFileName(originalFilename);
}
/**
* 删除
@ -265,21 +281,24 @@ public class MinioService {
* @param inputStream 输入流
* @param size 文件大小
* @param contentType 内容类型
* @return 实际写入的唯一对象名称
*/
public void uploadStream(String bucketName, String objectName, InputStream inputStream, long size, String contentType) {
public String uploadStream(String bucketName, String objectName, InputStream inputStream, long size, String contentType) {
String uniqueObjectName = buildUniqueObjectName(objectName);
try {
Map<String, String> extraHeaders = new HashMap<>();
extraHeaders.put("x-amz-acl", "public-read");
PutObjectArgs objectArgs = PutObjectArgs.builder()
.bucket(bucketName)
.object(objectName)
.object(uniqueObjectName)
.stream(inputStream, size, -1)
.extraHeaders(extraHeaders)
.contentType(contentType)
.build();
minioClient.putObject(objectArgs);
return uniqueObjectName;
} catch (Exception e) {
log.error("上传文件流失败, objectName: {}", objectName, e);
log.error("上传文件流失败, objectName: {}", uniqueObjectName, e);
throw new RuntimeException("上传文件流失败", e);
}
}

View File

@ -0,0 +1,45 @@
package com.cmvr.common.core.minio;
import org.junit.Test;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.IntStream;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
public class MinioServiceTest {
@Test
public void generatesUniqueNamesUnderConcurrency() {
int count = 10_000;
Set<String> names = ConcurrentHashMap.newKeySet();
IntStream.range(0, count).parallel()
.mapToObj(index -> MinioService.buildUniqueFileName("recording.mp4"))
.forEach(names::add);
assertEquals(count, names.size());
assertTrue(names.stream().allMatch(name ->
name.matches("\\d{13}-[0-9a-f]{32}\\.mp4")));
}
@Test
public void retainsExtensionAndSupportsExtensionlessNames() {
String imageName = MinioService.buildUniqueFileName("folder/archive/photo.JPG");
String extensionlessName = MinioService.buildUniqueFileName("README");
assertTrue(imageName.endsWith(".JPG"));
assertFalse(extensionlessName.contains("."));
}
@Test
public void retainsObjectDirectoryWhenGeneratingStreamUploadName() {
String objectName = MinioService.buildUniqueObjectName("inspection/alerts/source.jpg");
assertTrue(objectName.matches(
"inspection/alerts/\\d{13}-[0-9a-f]{32}\\.jpg"));
}
}

View File

@ -293,10 +293,11 @@ public class InspectionDetectionAlertServiceImpl implements IInspectionDetection
String extension = extensionFor(image.getMediaType());
String datePath = LocalDate.now().format(DateTimeFormatter.BASIC_ISO_DATE);
String objectName = IMAGE_OBJECT_PREFIX + "/" + datePath + "/" + alertId + extension;
String requestedObjectName = IMAGE_OBJECT_PREFIX + "/" + datePath + "/" + alertId + extension;
try (ByteArrayInputStream inputStream = new ByteArrayInputStream(imageBytes))
{
minioService.uploadStream(minioProperties.getBucketName(), objectName, inputStream,
String objectName = minioService.uploadStream(
minioProperties.getBucketName(), requestedObjectName, inputStream,
imageBytes.length, image.getMediaType());
return new StoredImage(objectName, buildMinioUrl(objectName));
}

View File

@ -90,6 +90,7 @@ public enum ActionEnum {
GET_CURRENT_PAGE("LLM", "GET_CURRENT_PAGE", "获取当前页面名称"),
AUDIO_EVENT_CLASSIFY("LLM", "AUDIO_EVENT_CLASSIFY", "声音类型检测"),
VIDEO_ANALYZE("LLM", "VIDEO_ANALYZE", "视频智能分析"),
IMAGE_ANALYZE("LLM", "IMAGE_ANALYZE", "通用图片分析"),
// 触控交互
TI_PATH_SEARCH("EDGE", "TI_PATH_SEARCH", "路径搜索"),

View File

@ -178,6 +178,10 @@ public class FlowModelBuilder {
// 设置动作类型
String actionStr = props.getString("action");
if (isLegacyImageAnalysisNode(nodeJson, props)) {
actionStr = ActionEnum.IMAGE_ANALYZE.getAction();
props.put("action", actionStr);
}
ActionEnum action = ActionEnum.fromAction(actionStr);
if (ObjUtil.isEmpty(action)) {
throw new IllegalArgumentException("无法识别的action:" + actionStr);
@ -205,6 +209,28 @@ public class FlowModelBuilder {
return wrapper;
}
/**
* Early image-analysis nodes could retain TOUCH_COORDINATES metadata while already
* carrying the new component type and profile. Normalize those persisted graphs at runtime.
*/
private static boolean isLegacyImageAnalysisNode(JSONObject nodeJson, JSONObject props) {
if (!"imageAnalysis".equalsIgnoreCase(nodeJson.getString("type"))) {
return false;
}
JSONArray params = props.getJSONArray("nodeParams");
if (params == null) {
return false;
}
for (int i = 0; i < params.size(); i++) {
JSONObject param = params.getJSONObject(i);
if ("profileCode".equals(param.getString("name"))
&& "common.image_analysis.v1".equals(param.getString("input"))) {
return true;
}
}
return false;
}
/**
* 解析分支条件逻辑
*/

View File

@ -1,5 +1,7 @@
package com.cmvr.test.flow.runtime.operator;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.JSONArray;
import com.alibaba.fastjson2.JSONObject;
import org.apache.commons.lang3.StringUtils;
@ -21,6 +23,21 @@ public final class FlowMediaParamResolver {
public static List<String> urls(JSONObject params, String name) {
Object value = params == null ? null : params.get(name);
if (value instanceof List<?> values) {
return normalizeUrls(values);
}
String text = normalizeUrl(value);
if (text != null && text.startsWith("[") && text.endsWith("]")) {
try {
JSONArray values = JSON.parseArray(text);
return normalizeUrls(values);
} catch (Exception ignored) {
// Keep treating malformed JSON as a scalar so validation can report it upstream.
}
}
return text == null ? Collections.emptyList() : List.of(text);
}
private static List<String> normalizeUrls(List<?> values) {
List<String> urls = new ArrayList<>();
for (Object item : values) {
String url = normalizeUrl(item);
@ -30,9 +47,6 @@ public final class FlowMediaParamResolver {
}
return urls;
}
String url = normalizeUrl(value);
return url == null ? Collections.emptyList() : List.of(url);
}
private static String normalizeUrl(Object value) {
return value == null ? null : StringUtils.trimToNull(String.valueOf(value));

View File

@ -11,6 +11,7 @@ import lombok.RequiredArgsConstructor;
import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.Set;
import java.util.UUID;
@ -19,27 +20,33 @@ import java.util.UUID;
public class LLMMediaAnalysisOperateService implements LLMOperateService {
private static final Set<String> VIDEO_ANALYSIS_MODES = Set.of("AUTO", "FAST", "ACCURATE");
private final MediaAnalysisClient mediaAnalysisClient;
@Override
public boolean supports(ActionEnum action) {
return action == ActionEnum.AUDIO_EVENT_CLASSIFY || action == ActionEnum.VIDEO_ANALYZE;
return action == ActionEnum.AUDIO_EVENT_CLASSIFY
|| action == ActionEnum.VIDEO_ANALYZE
|| action == ActionEnum.IMAGE_ANALYZE;
}
@Override
public TaskNodeExecuteResult execute(TaskNodeExecuteMessage message) {
JSONObject input = message.getInputParams();
boolean audio = message.getAction() == ActionEnum.AUDIO_EVENT_CLASSIFY;
String mediaParam = audio ? "audioUrl" : "videoUrl";
boolean image = message.getAction() == ActionEnum.IMAGE_ANALYZE;
String mediaParam = audio ? "audioUrl" : (image ? "imageUrl" : "videoUrl");
List<String> mediaUrls = FlowMediaParamResolver.urls(input, mediaParam);
String mediaUrl = FlowMediaParamResolver.lastUrl(input, mediaParam);
String profileCode = StringUtils.trimToNull(input.getString("profileCode"));
if (StringUtils.isBlank(mediaUrl)) {
throw new GlobalException("{}不能为空", audio ? "音频地址" : "视频地址");
throw new GlobalException("{}不能为空", audio ? "音频地址" : (image ? "图片地址" : "视频地址"));
}
if (profileCode == null) {
throw new GlobalException("分析场景不能为空");
}
if (image && mediaUrls.size() > 12) {
throw new GlobalException("图片分析一次最多支持12张图片");
}
JSONObject context = new JSONObject();
context.put("taskId", message.getTaskId());
@ -58,7 +65,7 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService {
options.put("expectedLabel", expectedLabel);
options.put("expectedLabelName", StringUtils.defaultIfBlank(expectedLabelName, expectedLabel));
}
} else {
} else if (!image) {
options.put("instruction", StringUtils.defaultString(input.getString("instruction")));
String analysisMode = StringUtils.defaultIfBlank(input.getString("analysisMode"), "AUTO")
.trim().toUpperCase();
@ -67,17 +74,42 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService {
}
options.put("analysisMode", analysisMode);
options.put("decisionPolicy", "FAIL_CLOSED");
options.put("detailedOutput", Boolean.TRUE.equals(input.getBoolean("detailedOutput")));
JSONObject tuning = buildVideoTuning(input.getJSONObject("analysisTuning"));
if (!tuning.isEmpty()) {
options.put("tuning", tuning);
}
} else {
String analysisMode = StringUtils.defaultIfBlank(input.getString("analysisMode"), "AUTO")
.trim().toUpperCase();
if (!VIDEO_ANALYSIS_MODES.contains(analysisMode)) {
throw new GlobalException("图片分析模式无效: {}", analysisMode);
}
String prompt = resolveImagePrompt(input);
if (StringUtils.isBlank(prompt)) {
throw new GlobalException("图片分析提示词不能为空");
}
options.put("prompt", prompt);
options.put("analysisMode", analysisMode);
String referenceImageUrl = FlowMediaParamResolver.lastUrl(input, "referenceImageUrl");
if (StringUtils.isNotBlank(referenceImageUrl)) {
options.put("referenceImageUrl", referenceImageUrl);
}
JSONObject tuning = buildImageTuning(input.getJSONObject("analysisTuning"));
if (!tuning.isEmpty()) {
options.put("tuning", tuning);
}
}
JSONObject request = new JSONObject();
request.put("requestId", buildRequestId(message));
request.put("analysisType", audio ? "AUDIO_CLASSIFICATION" : "VIDEO_ANALYSIS");
request.put("analysisType", audio ? "AUDIO_CLASSIFICATION"
: (image ? "IMAGE_ANALYSIS" : "VIDEO_ANALYSIS"));
request.put("profileCode", profileCode);
request.put("mediaUrl", mediaUrl);
if (image) {
request.put("mediaUrls", mediaUrls);
}
request.put("options", options);
request.put("context", context);
@ -94,6 +126,9 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService {
output.put("analysisModel", response.getJSONObject("model"));
output.put("analysisTimingMs", response.getLong("timingMs"));
output.put("mediaUrl", mediaUrl);
if (image) {
output.put("mediaUrls", mediaUrls);
}
if (audio) {
applyAudioVerdict(output, expectedLabel, expectedLabelName);
}
@ -127,20 +162,72 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService {
if (source == null || source.isEmpty()) {
return tuning;
}
Integer maxFrames = source.getInteger("maxFrames");
Double sampleFps = source.getDouble("sampleFps");
Integer maxWidth = source.getInteger("maxWidth");
Double confidenceThreshold = source.getDouble("confidenceThreshold");
Boolean fallbackToAccurate = source.getBoolean("fallbackToAccurate");
validateRange("抽帧数量", maxFrames, 2, 16);
validateRange("采样帧率", sampleFps, 0.25D, 6D);
validateRange("图片宽度", maxWidth, 640, 1280);
validateRange("通过置信度", confidenceThreshold, 0.5D, 0.95D);
if (maxFrames != null) tuning.put("maxFrames", maxFrames);
if (sampleFps != null) tuning.put("sampleFps", sampleFps);
if (maxWidth != null) tuning.put("maxWidth", maxWidth);
if (confidenceThreshold != null) tuning.put("confidenceThreshold", confidenceThreshold);
if (fallbackToAccurate != null) tuning.put("fallbackToAccurate", fallbackToAccurate);
return tuning;
}
private JSONObject buildImageTuning(JSONObject source) {
JSONObject tuning = new JSONObject();
if (source == null || source.isEmpty()) {
return tuning;
}
Integer maxWidth = source.getInteger("maxWidth");
Integer maxImages = source.getInteger("maxImages");
Integer maxOutputTokens = source.getInteger("maxOutputTokens");
validateRange("图片宽度", maxWidth, 640, 1600);
validateRange("图片数量", maxImages, 1, 12);
validateRange("最大输出长度", maxOutputTokens, 32, 4096);
if (maxWidth != null) tuning.put("maxWidth", maxWidth);
if (maxImages != null) tuning.put("maxImages", maxImages);
if (maxOutputTokens != null) tuning.put("maxOutputTokens", maxOutputTokens);
return tuning;
}
private String resolveImagePrompt(JSONObject input) {
String prompt = StringUtils.trimToNull(input.getString("prompt"));
if (prompt != null) {
return prompt;
}
// Old image-analysis nodes stored a fixed task type, target and pass criterion.
// Convert those fields into one neutral prompt so saved workflows remain runnable
// without retaining the old pass/fail semantics.
String target = StringUtils.trimToNull(input.getString("targetDescription"));
String instruction = StringUtils.trimToNull(input.getString("instruction"));
String taskType = StringUtils.trimToEmpty(input.getString("taskType")).toUpperCase();
StringBuilder legacyPrompt = new StringBuilder();
if (target != null) {
legacyPrompt.append(target);
}
if (instruction != null) {
if (!legacyPrompt.isEmpty()) {
legacyPrompt.append('\n');
}
legacyPrompt.append(instruction);
}
if ("ICON_TEMPLATE_LOCATE".equals(taskType)) {
if (!legacyPrompt.isEmpty()) {
legacyPrompt.append('\n');
}
legacyPrompt.append("最后一张图片是参考小图,请在其他图片中查找对应目标。");
}
if (!legacyPrompt.isEmpty()) {
legacyPrompt.append('\n');
}
legacyPrompt.append("只返回分析结果,不输出通过、未通过或置信度。");
return StringUtils.trimToEmpty(legacyPrompt.toString());
}
private void validateRange(String name, Number value, double minimum, double maximum) {
if (value != null && (value.doubleValue() < minimum || value.doubleValue() > maximum)) {
throw new GlobalException("{}必须在{}到{}之间", name, minimum, maximum);

View File

@ -327,7 +327,9 @@ public class FlowActionExecutorService {
private String executeLLMAction(ActionEnum action, FlowActionRequestVO req) {
log.info("LLM 执行动作: {}", action);
if (action == ActionEnum.AUDIO_EVENT_CLASSIFY || action == ActionEnum.VIDEO_ANALYZE) {
if (action == ActionEnum.AUDIO_EVENT_CLASSIFY
|| action == ActionEnum.VIDEO_ANALYZE
|| action == ActionEnum.IMAGE_ANALYZE) {
for (LLMOperateService service : llmOperateServices) {
if (service.supports(action)) {
return resultToString(service.execute(buildSingleNodeMessage(action, req.getPayload())));

View File

@ -0,0 +1,34 @@
package com.cmvr.test.flow.builder;
import com.cmvr.test.enums.ActionEnum;
import org.junit.Test;
import static org.junit.Assert.assertEquals;
public class FlowModelBuilderTest {
@Test
public void normalizesEarlyImageAnalysisNodeWithTouchAction() {
String flow = """
{
"nodes": [{
"id": "image-1",
"type": "imageAnalysis",
"properties": {
"name": "图片智能分析",
"action": "TOUCH_COORDINATES",
"nodeParams": [
{"name":"profileCode","type":"input","input":"common.image_analysis.v1"},
{"name":"imageUrl","type":"input","input":["http://minio/image.jpg"]}
]
}
}],
"edges": []
}
""";
FlowGraph graph = FlowModelBuilder.buildExecutableGraph(flow);
assertEquals(ActionEnum.IMAGE_ANALYZE, graph.getNodeMap().get("image-1").getAction());
}
}

View File

@ -12,6 +12,7 @@ import org.junit.Test;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
public class LLMMediaAnalysisOperateServiceTest {
@ -95,8 +96,9 @@ public class LLMMediaAnalysisOperateServiceTest {
.fluentPut("profileCode", "aima.power_video.v1")
.fluentPut("videoUrl", "https://files.example/video.mp4")
.fluentPut("analysisMode", "FAST")
.fluentPut("detailedOutput", true)
.fluentPut("analysisTuning", new JSONObject()
.fluentPut("maxFrames", 8)
.fluentPut("sampleFps", 2D)
.fluentPut("maxWidth", 960)
.fluentPut("confidenceThreshold", 0.8D)
.fluentPut("fallbackToAccurate", false))
@ -110,8 +112,10 @@ public class LLMMediaAnalysisOperateServiceTest {
assertEquals("FAST", captured.get().getJSONObject("options").getString("analysisMode"));
assertEquals("FAIL_CLOSED",
captured.get().getJSONObject("options").getString("decisionPolicy"));
assertEquals(true,
captured.get().getJSONObject("options").getBooleanValue("detailedOutput"));
JSONObject tuning = captured.get().getJSONObject("options").getJSONObject("tuning");
assertEquals(8, tuning.getIntValue("maxFrames"));
assertEquals(2D, tuning.getDoubleValue("sampleFps"), 0.0001D);
assertEquals(960, tuning.getIntValue("maxWidth"));
assertEquals(0.8D, tuning.getDoubleValue("confidenceThreshold"), 0.0001D);
assertEquals(false, tuning.getBooleanValue("fallbackToAccurate"));
@ -124,11 +128,85 @@ public class LLMMediaAnalysisOperateServiceTest {
.fluentPut("profileCode", "aima.power_video.v1")
.fluentPut("videoUrl", "https://files.example/video.mp4")
.fluentPut("analysisMode", "FAST")
.fluentPut("analysisTuning", new JSONObject().fluentPut("maxFrames", 30));
.fluentPut("analysisTuning", new JSONObject().fluentPut("sampleFps", 6.5D));
service.execute(message(ActionEnum.VIDEO_ANALYZE, input));
}
@Test
public void forwardsMultipleImagesAndReferenceImage() {
AtomicReference<JSONObject> captured = new AtomicReference<>();
MediaAnalysisClient client = request -> {
captured.set(request);
return new JSONObject()
.fluentPut("status", "SUCCEEDED")
.fluentPut("analysisType", "IMAGE_ANALYSIS")
.fluentPut("profileCode", "common.image_analysis.v1")
.fluentPut("result", new JSONObject()
.fluentPut("result", new JSONObject().fluentPut("x", 120).fluentPut("y", 80))
.fluentPut("resultText", "{\"x\":120,\"y\":80}"));
};
LLMMediaAnalysisOperateService service = new LLMMediaAnalysisOperateService(client);
JSONObject input = new JSONObject()
.fluentPut("profileCode", "common.image_analysis.v1")
.fluentPut("imageUrl", new JSONArray()
.fluentAdd("https://files.example/scene-1.jpg")
.fluentAdd("https://files.example/scene-2.jpg"))
.fluentPut("referenceImageUrl", "https://files.example/icon.jpg")
.fluentPut("prompt", "最后一张图片是参考图,输出匹配图标的中心坐标 JSON")
.fluentPut("analysisMode", "ACCURATE")
.fluentPut("analysisTuning", new JSONObject()
.fluentPut("maxWidth", 1280)
.fluentPut("maxImages", 8)
.fluentPut("maxOutputTokens", 256));
TaskNodeExecuteResult result = service.execute(message(ActionEnum.IMAGE_ANALYZE, input));
assertTrue(result.isSuccess());
assertEquals("IMAGE_ANALYSIS", captured.get().getString("analysisType"));
assertEquals(2, captured.get().getJSONArray("mediaUrls").size());
JSONObject options = captured.get().getJSONObject("options");
assertEquals("最后一张图片是参考图,输出匹配图标的中心坐标 JSON", options.getString("prompt"));
assertEquals("https://files.example/icon.jpg", options.getString("referenceImageUrl"));
assertEquals(8, options.getJSONObject("tuning").getIntValue("maxImages"));
assertEquals(256, options.getJSONObject("tuning").getIntValue("maxOutputTokens"));
assertEquals(2, result.getOutputParams().getJSONArray("mediaUrls").size());
assertNull(result.getOutputParams().get("passed"));
}
@Test
public void restoresStaticImageUrlArraySerializedByLegacyFlowParamModel() {
AtomicReference<JSONObject> captured = new AtomicReference<>();
MediaAnalysisClient client = request -> {
captured.set(request);
return new JSONObject()
.fluentPut("status", "SUCCEEDED")
.fluentPut("analysisType", "IMAGE_ANALYSIS")
.fluentPut("profileCode", "common.image_analysis.v1")
.fluentPut("result", new JSONObject()
.fluentPut("result", "图片正常")
.fluentPut("resultText", "图片正常"));
};
LLMMediaAnalysisOperateService service = new LLMMediaAnalysisOperateService(client);
JSONObject input = new JSONObject()
.fluentPut("profileCode", "common.image_analysis.v1")
.fluentPut("taskType", "GENERAL_INSPECTION")
.fluentPut("imageUrl", "[\"http://192.168.28.10:9000/cmvr-iot/path%2Fimage.jpg\"]")
.fluentPut("targetDescription", "检查图片")
.fluentPut("instruction", "满足要求则通过")
.fluentPut("analysisTuning", "{\"maxWidth\":1120,\"confidenceThreshold\":0.7,\"maxImages\":12}");
service.execute(message(ActionEnum.IMAGE_ANALYZE, input));
assertEquals("http://192.168.28.10:9000/cmvr-iot/path%2Fimage.jpg",
captured.get().getString("mediaUrl"));
assertEquals(1, captured.get().getJSONArray("mediaUrls").size());
assertEquals(1120, captured.get().getJSONObject("options")
.getJSONObject("tuning").getIntValue("maxWidth"));
assertTrue(captured.get().getJSONObject("options").getString("prompt").contains("检查图片"));
assertTrue(captured.get().getJSONObject("options").getString("prompt").contains("满足要求则通过"));
}
private TaskNodeExecuteMessage message(ActionEnum action, JSONObject input) {
TaskNodeExecuteMessage message = new TaskNodeExecuteMessage();
message.setAction(action);

View File

@ -153,8 +153,12 @@ public class CorpusInfoServiceImpl implements ICorpusInfoService {
}
try {
String objectName = "corpus/" + DateUtils.datePath() + "/" + corpusId + "_" + System.currentTimeMillis() + ".wav";
minioService.upload("cmvr-iot", "corpus/" + DateUtils.datePath(), file);
String uploadPath = "corpus/" + DateUtils.datePath();
String uploadedObjectName = minioService.upload("cmvr-iot", uploadPath, file);
if (StringUtils.isEmpty(uploadedObjectName)) {
throw new RuntimeException("上传音频失败");
}
String objectName = uploadPath + "/" + uploadedObjectName;
corpus.setAudioFilePath(objectName);
corpus.setFileSize(file.getSize());

View File

@ -285,11 +285,12 @@ public class TtsSynthesizeTaskServiceImpl implements ITtsSynthesizeTaskService {
if (audioBytes != null && audioBytes.length > 0) {
// 上传音频到Minio
String objectName = "tts/" + DateUtils.datePath() + "/" + task.getTaskCode() + ".wav";
String requestedObjectName = "tts/" + DateUtils.datePath() + "/" + task.getTaskCode() + ".wav";
ByteArrayInputStream inputStream = new ByteArrayInputStream(audioBytes);
// 使用MinioService上传
minioService.uploadStream(bucketName, objectName, inputStream, audioBytes.length, "audio/wav");
String objectName = minioService.uploadStream(
bucketName, requestedObjectName, inputStream, audioBytes.length, "audio/wav");
// 更新任务状态
task.setStatus("2");