feat(test): 完善测试任务控制与AI测试日志记录功能

- 在AimaTestLog中添加caseName字段并实现获取方法
- 扩展grpc协议定义包括CommandReasonCode和CommandExecutionState枚举
- 在系统命令中添加安全状态相关字段和操作接口
- 将AGV导航命令改为同步执行模式以提高可靠性
- 添加摄像头和麦克风服务的录像录音中止功能
- 实现测试任务终止时先清理媒体会话再停止设备的安全流程
- 为LLM媒体分析操作添加视频分析模式验证和调优参数支持
- 增加视频分析参数范围校验和测试用例覆盖
This commit is contained in:
lixiaolong 2026-08-17 16:58:31 +08:00
parent 861252f066
commit d3498e0178
21 changed files with 21538 additions and 145 deletions

View File

@ -151,13 +151,15 @@ public class AimaFlowExecutionListener implements FlowExecutionListener {
String imageUrl = FlowMediaParamResolver.lastUrl(outputParams, "imageUrl");
String videoUrl = FlowMediaParamResolver.lastUrl(outputParams, "videoUrl");
String testCaseId = resolveTestCaseId(instId, taskInstance.getTaskId(), event);
// 获取casename
String casename = aimaTaskService.getCaseName(testCaseId);
// 保存日志到数据库
AimaTestLog testLog = AimaTestLog.builder()
.id(logId)
.taskInstanceId(taskInstance.getId())
.taskId(taskInstance.getTaskId())
.testCaseId(testCaseId)
.caseName(casename)
.itemId(itemId)
.itemOrder(event.getItemOrder())
.itemOccurrence(event.getItemOccurrence())

View File

@ -36,4 +36,6 @@ public interface IAimaTaskService {
* @return 测试用例ID列表
*/
List<AimaTestCase> getTestCaseIds(String taskId);
String getCaseName(String testCaseId);
}

View File

@ -204,4 +204,11 @@ public class AimaTaskServiceImpl extends ServiceImpl<AimaTaskMapper, AimaTask> i
public List<AimaTestCase> getTestCaseIds(String taskId) {
return aimaTaskTestCaseMapper.selectTestCasesByTaskId(taskId);
}
@Override
public String getCaseName(String testCaseId) {
AimaTestCase aimaTestCase = aimaTestCaseMapper.selectAimaTestCaseById(testCaseId);
return aimaTestCase != null ? aimaTestCase.getCaseName() : null;
}
}

View File

@ -56,6 +56,9 @@ public interface EdgeCameraService {
/** Stops and finalizes a backend-local workflow recording. */
File stopWorkflowRecording(EdgeCommonVO edgeCommonVO, String recordingId);
/** Aborts every backend-local workflow recording owned by the robot. */
void abortWorkflowRecordings(String robotId);
/**
* 获取相机图片组
*/

View File

@ -33,6 +33,9 @@ public interface EdgeMicrophoneService {
/** Stops and finalizes a backend-local workflow recording. */
File stopWorkflowRecording(EdgeCommonVO edgeCommonVO, String recordingId);
/** Aborts every backend-local workflow recording owned by the robot. */
void abortWorkflowRecordings(String robotId);
/**
* 暂停录音
*/

View File

@ -91,7 +91,7 @@ public class EdgeAgvServiceImpl implements EdgeAgvService {
AgvCommand.AgvNavigateToStationCommand.Request request = AgvCommand.AgvNavigateToStationCommand.Request.newBuilder()
.setHeader(EdgeCommonUtil.buildRequest(edgeCommonVO.getDeviceId()))
.setOptions(AgvUtils.AgvMotionOptions.newBuilder()
.setAsynchronous(true)
.setAsynchronous(false)
.build())
.setStationId(stationId)
.build();

View File

@ -254,6 +254,18 @@ public class EdgeCameraServiceImpl implements EdgeCameraService, EdgeStreamServi
});
}
@Override
public void abortWorkflowRecordings(String robotId) {
workflowRecordings.forEach((recordingId, session) -> {
if (StrUtil.equals(robotId, session.robotId)
&& workflowRecordings.remove(recordingId, session)) {
log.info("终止机器人关联的工作流录像: recordingId={}, robotId={}, deviceId={}",
recordingId, session.robotId, session.deviceId);
session.abort();
}
});
}
@Scheduled(fixedDelay = 60_000L)
public void cleanExpiredWorkflowRecordings() {
workflowRecordings.forEach((recordingId, session) -> {

View File

@ -184,6 +184,22 @@ public class EdgeMicrophoneServiceImpl implements EdgeMicrophoneService {
});
}
@Override
public void abortWorkflowRecordings(String robotId) {
workflowRecordings.forEach((recordingId, session) -> {
if (StrUtil.equals(robotId, session.robotId)
&& workflowRecordings.remove(recordingId, session)) {
workflowRecordingTargets.remove(session.targetKey(), recordingId);
logger.info("终止机器人关联的工作流录音: recordingId={}, robotId={}, deviceId={}",
recordingId, session.robotId, session.deviceId);
session.abort();
}
});
// 会话创建失败或异常退出时也不能遗留目标占用记录。
String targetPrefix = String.valueOf(robotId) + '\0';
workflowRecordingTargets.keySet().removeIf(key -> key.startsWith(targetPrefix));
}
@Scheduled(fixedDelay = 60_000L)
public void cleanExpiredWorkflowRecordings() {
workflowRecordings.forEach((recordingId, session) -> {

View File

@ -201,6 +201,68 @@ public final class SystemServiceGrpc {
return getExecuteActionQueueMethod;
}
private static volatile io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.GetSafetyStateCommand.Request,
cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback> getGetSafetyStateMethod;
@io.grpc.stub.annotations.RpcMethod(
fullMethodName = SERVICE_NAME + '/' + "GetSafetyState",
requestType = cmvr.api.SafetyCommand.GetSafetyStateCommand.Request.class,
responseType = cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback.class,
methodType = io.grpc.MethodDescriptor.MethodType.UNARY)
public static io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.GetSafetyStateCommand.Request,
cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback> getGetSafetyStateMethod() {
io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.GetSafetyStateCommand.Request, cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback> getGetSafetyStateMethod;
if ((getGetSafetyStateMethod = SystemServiceGrpc.getGetSafetyStateMethod) == null) {
synchronized (SystemServiceGrpc.class) {
if ((getGetSafetyStateMethod = SystemServiceGrpc.getGetSafetyStateMethod) == null) {
SystemServiceGrpc.getGetSafetyStateMethod = getGetSafetyStateMethod =
io.grpc.MethodDescriptor.<cmvr.api.SafetyCommand.GetSafetyStateCommand.Request, cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback>newBuilder()
.setType(io.grpc.MethodDescriptor.MethodType.UNARY)
.setFullMethodName(generateFullMethodName(SERVICE_NAME, "GetSafetyState"))
.setSampledToLocalTracing(true)
.setRequestMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
cmvr.api.SafetyCommand.GetSafetyStateCommand.Request.getDefaultInstance()))
.setResponseMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback.getDefaultInstance()))
.setSchemaDescriptor(new SystemServiceMethodDescriptorSupplier("GetSafetyState"))
.build();
}
}
}
return getGetSafetyStateMethod;
}
private static volatile io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request,
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback> getRecoverSafetyStateMethod;
@io.grpc.stub.annotations.RpcMethod(
fullMethodName = SERVICE_NAME + '/' + "RecoverSafetyState",
requestType = cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request.class,
responseType = cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback.class,
methodType = io.grpc.MethodDescriptor.MethodType.UNARY)
public static io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request,
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback> getRecoverSafetyStateMethod() {
io.grpc.MethodDescriptor<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request, cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback> getRecoverSafetyStateMethod;
if ((getRecoverSafetyStateMethod = SystemServiceGrpc.getRecoverSafetyStateMethod) == null) {
synchronized (SystemServiceGrpc.class) {
if ((getRecoverSafetyStateMethod = SystemServiceGrpc.getRecoverSafetyStateMethod) == null) {
SystemServiceGrpc.getRecoverSafetyStateMethod = getRecoverSafetyStateMethod =
io.grpc.MethodDescriptor.<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request, cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback>newBuilder()
.setType(io.grpc.MethodDescriptor.MethodType.UNARY)
.setFullMethodName(generateFullMethodName(SERVICE_NAME, "RecoverSafetyState"))
.setSampledToLocalTracing(true)
.setRequestMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request.getDefaultInstance()))
.setResponseMarshaller(io.grpc.protobuf.ProtoUtils.marshaller(
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback.getDefaultInstance()))
.setSchemaDescriptor(new SystemServiceMethodDescriptorSupplier("RecoverSafetyState"))
.build();
}
}
}
return getRecoverSafetyStateMethod;
}
/**
* Creates a new async stub that supports all call types for the service
*/
@ -291,6 +353,20 @@ public final class SystemServiceGrpc {
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getExecuteActionQueueMethod(), responseObserver);
}
/**
*/
public void getSafetyState(cmvr.api.SafetyCommand.GetSafetyStateCommand.Request request,
io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback> responseObserver) {
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getGetSafetyStateMethod(), responseObserver);
}
/**
*/
public void recoverSafetyState(cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request request,
io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback> responseObserver) {
io.grpc.stub.ServerCalls.asyncUnimplementedUnaryCall(getRecoverSafetyStateMethod(), responseObserver);
}
@java.lang.Override public final io.grpc.ServerServiceDefinition bindService() {
return io.grpc.ServerServiceDefinition.builder(getServiceDescriptor())
.addMethod(
@ -335,6 +411,20 @@ public final class SystemServiceGrpc {
cmvr.api.SystemCommand.ActionQueueCommand.Request,
cmvr.api.SystemCommand.ActionQueueCommand.Feedback>(
this, METHODID_EXECUTE_ACTION_QUEUE)))
.addMethod(
getGetSafetyStateMethod(),
io.grpc.stub.ServerCalls.asyncUnaryCall(
new MethodHandlers<
cmvr.api.SafetyCommand.GetSafetyStateCommand.Request,
cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback>(
this, METHODID_GET_SAFETY_STATE)))
.addMethod(
getRecoverSafetyStateMethod(),
io.grpc.stub.ServerCalls.asyncUnaryCall(
new MethodHandlers<
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request,
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback>(
this, METHODID_RECOVER_SAFETY_STATE)))
.build();
}
}
@ -400,6 +490,22 @@ public final class SystemServiceGrpc {
io.grpc.stub.ClientCalls.asyncUnaryCall(
getChannel().newCall(getExecuteActionQueueMethod(), getCallOptions()), request, responseObserver);
}
/**
*/
public void getSafetyState(cmvr.api.SafetyCommand.GetSafetyStateCommand.Request request,
io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback> responseObserver) {
io.grpc.stub.ClientCalls.asyncUnaryCall(
getChannel().newCall(getGetSafetyStateMethod(), getCallOptions()), request, responseObserver);
}
/**
*/
public void recoverSafetyState(cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request request,
io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback> responseObserver) {
io.grpc.stub.ClientCalls.asyncUnaryCall(
getChannel().newCall(getRecoverSafetyStateMethod(), getCallOptions()), request, responseObserver);
}
}
/**
@ -457,6 +563,20 @@ public final class SystemServiceGrpc {
return io.grpc.stub.ClientCalls.blockingUnaryCall(
getChannel(), getExecuteActionQueueMethod(), getCallOptions(), request);
}
/**
*/
public cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback getSafetyState(cmvr.api.SafetyCommand.GetSafetyStateCommand.Request request) {
return io.grpc.stub.ClientCalls.blockingUnaryCall(
getChannel(), getGetSafetyStateMethod(), getCallOptions(), request);
}
/**
*/
public cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback recoverSafetyState(cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request request) {
return io.grpc.stub.ClientCalls.blockingUnaryCall(
getChannel(), getRecoverSafetyStateMethod(), getCallOptions(), request);
}
}
/**
@ -520,6 +640,22 @@ public final class SystemServiceGrpc {
return io.grpc.stub.ClientCalls.futureUnaryCall(
getChannel().newCall(getExecuteActionQueueMethod(), getCallOptions()), request);
}
/**
*/
public com.google.common.util.concurrent.ListenableFuture<cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback> getSafetyState(
cmvr.api.SafetyCommand.GetSafetyStateCommand.Request request) {
return io.grpc.stub.ClientCalls.futureUnaryCall(
getChannel().newCall(getGetSafetyStateMethod(), getCallOptions()), request);
}
/**
*/
public com.google.common.util.concurrent.ListenableFuture<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback> recoverSafetyState(
cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request request) {
return io.grpc.stub.ClientCalls.futureUnaryCall(
getChannel().newCall(getRecoverSafetyStateMethod(), getCallOptions()), request);
}
}
private static final int METHODID_GET_SYSTEM_INFO = 0;
@ -528,6 +664,8 @@ public final class SystemServiceGrpc {
private static final int METHODID_UPDATE_PARAMS = 3;
private static final int METHODID_STOP_ALL = 4;
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 class MethodHandlers<Req, Resp> implements
io.grpc.stub.ServerCalls.UnaryMethod<Req, Resp>,
@ -570,6 +708,14 @@ public final class SystemServiceGrpc {
serviceImpl.executeActionQueue((cmvr.api.SystemCommand.ActionQueueCommand.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.SystemCommand.ActionQueueCommand.Feedback>) responseObserver);
break;
case METHODID_GET_SAFETY_STATE:
serviceImpl.getSafetyState((cmvr.api.SafetyCommand.GetSafetyStateCommand.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.GetSafetyStateCommand.Feedback>) responseObserver);
break;
case METHODID_RECOVER_SAFETY_STATE:
serviceImpl.recoverSafetyState((cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Request) request,
(io.grpc.stub.StreamObserver<cmvr.api.SafetyCommand.RecoverSafetyStateCommand.Feedback>) responseObserver);
break;
default:
throw new AssertionError();
}
@ -637,6 +783,8 @@ public final class SystemServiceGrpc {
.addMethod(getUpdateParamsMethod())
.addMethod(getStopAllMethod())
.addMethod(getExecuteActionQueueMethod())
.addMethod(getGetSafetyStateMethod())
.addMethod(getRecoverSafetyStateMethod())
.build();
}
}

View File

@ -24,30 +24,38 @@ public final class SystemServiceOuterClass {
static {
java.lang.String[] descriptorData = {
"\n\035cmvr/api/system_service.proto\022\010cmvr.ap" +
"i\032\035cmvr/api/system_command.proto2\331\004\n\rSys" +
"temService\022b\n\rGetSystemInfo\022&.cmvr.api.G" +
"etSystemInfoCommand.Request\032\'.cmvr.api.G" +
"etSystemInfoCommand.Feedback\"\000\022h\n\017GetSys" +
"temStatus\022(.cmvr.api.GetSystemStatusComm" +
"and.Request\032).cmvr.api.GetSystemStatusCo" +
"mmand.Feedback\"\000\022b\n\rGetDeviceList\022&.cmvr" +
".api.GetDeviceListCommand.Request\032\'.cmvr" +
".api.GetDeviceListCommand.Feedback\"\000\022_\n\014" +
"UpdateParams\022%.cmvr.api.UpdateParamsComm" +
"and.Request\032&.cmvr.api.UpdateParamsComma" +
"nd.Feedback\"\000\022P\n\007StopAll\022 .cmvr.api.Stop" +
"AllCommand.Request\032!.cmvr.api.StopAllCom" +
"mand.Feedback\"\000\022c\n\022ExecuteActionQueue\022$." +
"cmvr.api.ActionQueueCommand.Request\032%.cm" +
"vr.api.ActionQueueCommand.Feedback\"\000b\006pr" +
"oto3"
"i\032\035cmvr/api/system_command.proto\032\035cmvr/a" +
"pi/safety_command.proto2\263\006\n\rSystemServic" +
"e\022b\n\rGetSystemInfo\022&.cmvr.api.GetSystemI" +
"nfoCommand.Request\032\'.cmvr.api.GetSystemI" +
"nfoCommand.Feedback\"\000\022h\n\017GetSystemStatus" +
"\022(.cmvr.api.GetSystemStatusCommand.Reque" +
"st\032).cmvr.api.GetSystemStatusCommand.Fee" +
"dback\"\000\022b\n\rGetDeviceList\022&.cmvr.api.GetD" +
"eviceListCommand.Request\032\'.cmvr.api.GetD" +
"eviceListCommand.Feedback\"\000\022_\n\014UpdatePar" +
"ams\022%.cmvr.api.UpdateParamsCommand.Reque" +
"st\032&.cmvr.api.UpdateParamsCommand.Feedba" +
"ck\"\000\022P\n\007StopAll\022 .cmvr.api.StopAllComman" +
"d.Request\032!.cmvr.api.StopAllCommand.Feed" +
"back\"\000\022c\n\022ExecuteActionQueue\022$.cmvr.api." +
"ActionQueueCommand.Request\032%.cmvr.api.Ac" +
"tionQueueCommand.Feedback\"\000\022e\n\016GetSafety" +
"State\022\'.cmvr.api.GetSafetyStateCommand.R" +
"equest\032(.cmvr.api.GetSafetyStateCommand." +
"Feedback\"\000\022q\n\022RecoverSafetyState\022+.cmvr." +
"api.RecoverSafetyStateCommand.Request\032,." +
"cmvr.api.RecoverSafetyStateCommand.Feedb" +
"ack\"\000b\006proto3"
};
descriptor = com.google.protobuf.Descriptors.FileDescriptor
.internalBuildGeneratedFileFrom(descriptorData,
new com.google.protobuf.Descriptors.FileDescriptor[] {
cmvr.api.SystemCommand.getDescriptor(),
cmvr.api.SafetyCommand.getDescriptor(),
});
cmvr.api.SystemCommand.getDescriptor();
cmvr.api.SafetyCommand.getDescriptor();
}
// @@protoc_insertion_point(outer_class_scope)

View File

@ -4,6 +4,59 @@ package cmvr.api;
import "google/protobuf/timestamp.proto";
enum CommandReasonCode {
COMMAND_REASON_CODE_UNSPECIFIED = 0;
COMMAND_REASON_CODE_NONE = 1;
COMMAND_REASON_CODE_INVALID_ARGUMENT = 2;
COMMAND_REASON_CODE_UNAUTHENTICATED = 3;
COMMAND_REASON_CODE_PERMISSION_DENIED = 4;
COMMAND_REASON_CODE_RECOVERY_RPC_DISABLED = 5;
COMMAND_REASON_CODE_DEVICE_NOT_FOUND = 6;
COMMAND_REASON_CODE_DEVICE_UNAVAILABLE = 7;
COMMAND_REASON_CODE_UNSUPPORTED_COMMAND = 8;
COMMAND_REASON_CODE_SYSTEM_STARTING = 9;
COMMAND_REASON_CODE_SYSTEM_STOPPING = 10;
COMMAND_REASON_CODE_SAFETY_LATCHED = 11;
COMMAND_REASON_CODE_SAFETY_STATE_MISSING = 12;
COMMAND_REASON_CODE_SAFETY_STATE_STALE = 13;
COMMAND_REASON_CODE_HARDWARE_UNSAFE = 14;
COMMAND_REASON_CODE_EMERGENCY_STOP_ACTIVE = 15;
COMMAND_REASON_CODE_PROTECTIVE_STOP_ACTIVE = 16;
COMMAND_REASON_CODE_DEVICE_DISCONNECTED = 17;
COMMAND_REASON_CODE_DEVICE_FAULT = 18;
COMMAND_REASON_CODE_DEVICE_NOT_READY = 19;
COMMAND_REASON_CODE_DEVICE_STILL_MOVING = 20;
COMMAND_REASON_CODE_CONTROL_BUSY = 21;
COMMAND_REASON_CODE_GENERATION_MISMATCH = 22;
COMMAND_REASON_CODE_COMMAND_ID_REQUIRED = 23;
COMMAND_REASON_CODE_COMMAND_ID_CONFLICT = 24;
COMMAND_REASON_CODE_RESULT_EVICTED = 25;
COMMAND_REASON_CODE_LEDGER_EXHAUSTED = 26;
COMMAND_REASON_CODE_BACKPRESSURE = 27;
COMMAND_REASON_CODE_DEADLINE_EXCEEDED_BEFORE_DISPATCH = 28;
COMMAND_REASON_CODE_OUTCOME_UNKNOWN = 29;
COMMAND_REASON_CODE_PARTICIPANT_TIMEOUT = 30;
COMMAND_REASON_CODE_STOP_UNCONFIRMED = 31;
COMMAND_REASON_CODE_RECOVERY_EPOCH_MISMATCH = 32;
COMMAND_REASON_CODE_RECOVERY_REASON_REQUIRED = 33;
COMMAND_REASON_CODE_RECOVERY_AUDIT_FAILED = 34;
COMMAND_REASON_CODE_INTERNAL_ERROR = 35;
}
enum CommandExecutionState {
COMMAND_EXECUTION_STATE_UNSPECIFIED = 0;
COMMAND_EXECUTION_STATE_RECEIVED = 1;
COMMAND_EXECUTION_STATE_RESERVED = 2;
COMMAND_EXECUTION_STATE_REJECTED_BEFORE_DISPATCH = 3;
COMMAND_EXECUTION_STATE_ADMITTED = 4;
COMMAND_EXECUTION_STATE_DISPATCHING = 5;
COMMAND_EXECUTION_STATE_ACCEPTED_BY_HARDWARE = 6;
COMMAND_EXECUTION_STATE_COMPLETED = 7;
COMMAND_EXECUTION_STATE_FAILED = 8;
COMMAND_EXECUTION_STATE_CANCELED_BEFORE_DISPATCH = 9;
COMMAND_EXECUTION_STATE_OUTCOME_UNKNOWN = 10;
}
message DeviceLifecycle {
enum Lifecycle {
STATE_INIT = 0;
@ -20,12 +73,22 @@ message CommandHeader {
message Request {
string device_id = 1; // 目标设备名称
google.protobuf.Timestamp timestamp = 2; // 请求时间
string command_id = 3;
string expected_service_instance_id = 4;
optional uint64 expected_device_generation = 5;
uint32 valid_for_ms = 6;
}
message Feedback {
bool success = 1; // 是否成功
string error_message = 2; // 错误信息(成功时为空)
google.protobuf.Timestamp timestamp = 3; // 回复时间
CommandReasonCode reason_code = 4;
string command_id = 5;
string service_instance_id = 6;
uint64 safety_epoch = 7;
uint64 device_generation = 8;
CommandExecutionState execution_state = 9;
}
}

View File

@ -0,0 +1,186 @@
syntax = "proto3";
package cmvr.api;
import "cmvr/api/common.proto";
enum SafetyTriState {
SAFETY_TRI_STATE_UNKNOWN = 0;
SAFETY_TRI_STATE_FALSE = 1;
SAFETY_TRI_STATE_TRUE = 2;
}
enum SafetyCondition {
SAFETY_CONDITION_UNKNOWN = 0;
SAFETY_CONDITION_NOMINAL = 1;
SAFETY_CONDITION_RESTRICTED = 2;
SAFETY_CONDITION_UNSAFE = 3;
}
enum SystemAdmissionState {
SYSTEM_ADMISSION_STATE_UNSPECIFIED = 0;
SYSTEM_ADMISSION_STATE_STARTING = 1;
SYSTEM_ADMISSION_STATE_OPEN = 2;
SYSTEM_ADMISSION_STATE_STOPPING = 3;
SYSTEM_ADMISSION_STATE_LATCHED = 4;
SYSTEM_ADMISSION_STATE_RECOVERING = 5;
SYSTEM_ADMISSION_STATE_SHUTTING_DOWN = 6;
}
enum DeviceAdmissionState {
DEVICE_ADMISSION_STATE_UNSPECIFIED = 0;
DEVICE_ADMISSION_STATE_OBSERVING = 1;
DEVICE_ADMISSION_STATE_OPEN = 2;
DEVICE_ADMISSION_STATE_BLOCKED = 3;
DEVICE_ADMISSION_STATE_QUARANTINED = 4;
DEVICE_ADMISSION_STATE_RECOVERING = 5;
DEVICE_ADMISSION_STATE_REMOVED = 6;
}
enum SafetyBlockerScope {
SAFETY_BLOCKER_SCOPE_UNSPECIFIED = 0;
SAFETY_BLOCKER_SCOPE_DEVICE = 1;
SAFETY_BLOCKER_SCOPE_SYSTEM = 2;
}
enum SafetyRecoveryRequirement {
SAFETY_RECOVERY_REQUIREMENT_UNSPECIFIED = 0;
SAFETY_RECOVERY_REQUIREMENT_REFRESH_ONLY = 1;
SAFETY_RECOVERY_REQUIREMENT_CLEAR_SOFTWARE_LATCH = 2;
SAFETY_RECOVERY_REQUIREMENT_HARDWARE_RELEASE_REQUIRED = 3;
SAFETY_RECOVERY_REQUIREMENT_MANUAL_INSPECTION_REQUIRED = 4;
}
enum SafetyOperationResult {
SAFETY_OPERATION_RESULT_UNSPECIFIED = 0;
SAFETY_OPERATION_RESULT_SUCCEEDED = 1;
SAFETY_OPERATION_RESULT_RECOVERED = 2;
SAFETY_OPERATION_RESULT_VERIFIED_BUT_STILL_BLOCKED = 3;
SAFETY_OPERATION_RESULT_BLOCKER_REMAINS = 4;
SAFETY_OPERATION_RESULT_EPOCH_MISMATCH = 5;
SAFETY_OPERATION_RESULT_NOTHING_TO_RECOVER = 6;
SAFETY_OPERATION_RESULT_TIMED_OUT = 7;
SAFETY_OPERATION_RESULT_FAILED = 8;
}
message DeviceIdList {
repeated string device_ids = 1;
}
message SafetyScope {
oneof target {
bool all_devices = 1;
DeviceIdList devices = 2;
}
}
message SafetyBlockerInfo {
CommandReasonCode reason_code = 1;
SafetyBlockerScope scope = 2;
SafetyRecoveryRequirement recovery_requirement = 3;
string source_id = 4;
string operation_id = 5;
uint64 first_observed_at_unix_ms = 6;
uint64 last_observed_at_unix_ms = 7;
}
message DeviceSafetyStateInfo {
string device_id = 1;
string device_kind = 2;
string policy_family = 3;
string lifecycle_state = 4;
string health_state = 5;
DeviceAdmissionState admission_state = 6;
SafetyCondition condition = 7;
bool has_sample = 8;
bool snapshot_fresh = 9;
uint64 sample_age_ms = 10;
uint64 sample_sequence = 11;
uint64 observed_at_unix_ms = 12;
uint64 device_generation = 13;
SafetyTriState connected = 14;
SafetyTriState operational_ready = 15;
SafetyTriState quiescent = 16;
SafetyTriState motion_active = 17;
SafetyTriState actuator_enabled = 18;
SafetyTriState emergency_stop_active = 19;
SafetyTriState protective_stop_active = 20;
SafetyTriState fault_active = 21;
repeated SafetyBlockerInfo blockers = 22;
}
message SafetyOperationTargetResult {
string target_id = 1;
SafetyOperationResult result = 2;
CommandReasonCode reason_code = 3;
string detail = 4;
DeviceAdmissionState before_state = 5;
DeviceAdmissionState after_state = 6;
}
message SafetyParticipantResultInfo {
bool recorded = 1;
bool success = 2;
CommandReasonCode reason_code = 3;
string detail = 4;
}
message SafetyParticipantStateInfo {
string participant_id = 1;
string phase = 2;
bool required = 3;
bool registered = 4;
bool barrier_active = 5;
bool barrier_retained = 6;
string operation_id = 7;
uint64 safety_epoch = 8;
SafetyParticipantResultInfo last_request = 9;
SafetyParticipantResultInfo last_verify = 10;
SafetyParticipantResultInfo last_release = 11;
}
message GetSafetyStateCommand {
message Request {
SafetyScope scope = 1;
}
message Feedback {
CommandHeader.Feedback header = 1;
SystemAdmissionState system_state = 2;
uint64 safety_epoch = 3;
string control_service_instance_id = 4;
string enforcement_mode = 5;
repeated DeviceSafetyStateInfo devices = 6;
string active_operation_id = 7;
string active_operation_phase = 8;
uint64 sampled_at_unix_ms = 9;
repeated SafetyParticipantStateInfo participants = 10;
}
}
message RecoverSafetyStateCommand {
enum Mode {
MODE_UNSPECIFIED = 0;
VERIFY_ONLY = 1;
CLEAR_SOFTWARE_LATCH = 2;
}
message Request {
string recovery_id = 1;
SafetyScope scope = 2;
uint64 expected_safety_epoch = 3;
Mode mode = 4;
string reason = 5;
uint32 timeout_ms = 6;
}
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

@ -3,6 +3,7 @@ syntax = "proto3";
import "cmvr/api/agv_command.proto";
import "cmvr/api/arm_command.proto";
import "cmvr/api/common.proto";
import "cmvr/api/safety_command.proto";
package cmvr.api;
@ -109,6 +110,16 @@ message GetSystemInfoCommand {
// Changes whenever the in-process ActionQueue idempotency ledger is
// recreated. Clients bind submissions and retries to this value.
string action_service_instance_id = 8;
// Effective server-side control-plane settings. These fields describe
// what is running, not merely what the configuration requested.
string grpc_transport_security = 9;
string grpc_authentication = 10;
string grpc_recovery_exposure = 11;
bool grpc_insecure_non_loopback = 12;
string control_service_instance_id = 13;
string safety_enforcement_mode = 14;
uint32 safety_schema_version = 15;
}
}
@ -141,10 +152,18 @@ message UpdateParamsCommand {
message StopAllCommand {
message Request {
CommandHeader.Request header = 1;
string operation_id = 2;
string expected_service_instance_id = 3;
uint32 timeout_ms = 4;
}
message Feedback {
CommandHeader.Feedback header = 1;
string operation_id = 2;
uint64 previous_safety_epoch = 3;
uint64 current_safety_epoch = 4;
SystemAdmissionState system_state = 5;
repeated SafetyOperationTargetResult targets = 6;
}
}

View File

@ -1,6 +1,7 @@
syntax = "proto3";
import "cmvr/api/system_command.proto";
import "cmvr/api/safety_command.proto";
package cmvr.api;
@ -15,4 +16,7 @@ service SystemService {
rpc StopAll(StopAllCommand.Request) returns (StopAllCommand.Feedback) {}
rpc ExecuteActionQueue(ActionQueueCommand.Request) returns (ActionQueueCommand.Feedback) {}
rpc GetSafetyState(GetSafetyStateCommand.Request) returns (GetSafetyStateCommand.Feedback) {}
rpc RecoverSafetyState(RecoverSafetyStateCommand.Request) returns (RecoverSafetyStateCommand.Feedback) {}
}

View File

@ -5,6 +5,8 @@ import cn.hutool.core.util.StrUtil;
import com.alibaba.fastjson2.JSON;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.edge.client.service.EdgeSystemService;
import com.cmvr.edge.client.service.EdgeCameraService;
import com.cmvr.edge.client.service.EdgeMicrophoneService;
import com.cmvr.test.enums.NodeTypeEnum;
import com.cmvr.test.enums.TaskStatusEnum;
import com.cmvr.test.enums.TerminalStatusEnum;
@ -54,6 +56,8 @@ public class FlowControlService {
private final TaskThreadRegistry taskThreadRegistry;
private final FlowItemExecutor flowItemExecutor;
private final EdgeSystemService edgeSystemService;
private final EdgeMicrophoneService edgeMicrophoneService;
private final EdgeCameraService edgeCameraService;
private final ITeTaskInstService taskInstService;
/**
@ -69,13 +73,12 @@ public class FlowControlService {
if (ctx.isStopped()) {
return;
}
// 异步终止
stopAllRobots(ctx);
// 统一记录日志 + 设置上下文状态 + 数据库状态
// 先阻止节点继续执行,再清理设备与本地媒体会话,避免终止过程中重新登记录音/录像。
taskInstHolder.syncStatus(instId, TaskStatusEnum.STOPPED);
taskInstHolder.unRegister(instId);
taskThreadRegistry.interruptAll(instId);
// 终止必须等设备和后端媒体会话完成清理后再返回,防止立即重跑命中旧录音/录像状态。
stopAllRobots(ctx);
taskInstHolder.unRegister(instId);
// 被中断的节点 状态修改成 PAUSED
for (TaskNodeExecuteContext pendingNode : taskThreadRegistry.getPendingNodes(instId)) {
String nodeId = pendingNode.getNode().getNodeId();
@ -144,14 +147,15 @@ public class FlowControlService {
persisted.getId(), exception.getMessage());
}
}
robotIds.forEach(robotId -> executor.execute(() -> {
robotIds.forEach(robotId -> {
abortLocalMediaRecordings(robotId);
try {
edgeSystemService.stopAll(robotId);
} catch (RuntimeException exception) {
log.warn("终止遗留任务设备失败: instId={}, robotId={}, error={}",
persisted.getId(), robotId, exception.getMessage());
}
}));
});
}
/**
@ -306,7 +310,14 @@ public class FlowControlService {
robotIds.stream()
.filter(robotId -> robotId != null && !robotId.isBlank())
.distinct()
.forEach(robotId -> executor.execute(() -> edgeSystemService.stopAll(robotId)));
.forEach(robotId -> {
abortLocalMediaRecordings(robotId);
try {
edgeSystemService.stopAll(robotId);
} catch (RuntimeException exception) {
log.warn("终止机器人设备失败: robotId={}, error={}", robotId, exception.getMessage());
}
});
}
private boolean stopAllRobotsAndWait(TaskContext context) {
@ -318,6 +329,7 @@ public class FlowControlService {
.filter(robotId -> robotId != null && !robotId.isBlank())
.distinct()
.collect(Collectors.toList())) {
abortLocalMediaRecordings(robotId);
try {
edgeSystemService.stopAll(robotId);
} catch (RuntimeException exception) {
@ -329,6 +341,19 @@ public class FlowControlService {
return success;
}
private void abortLocalMediaRecordings(String robotId) {
try {
edgeMicrophoneService.abortWorkflowRecordings(robotId);
} catch (RuntimeException exception) {
log.warn("清理工作流录音会话失败: robotId={}, error={}", robotId, exception.getMessage());
}
try {
edgeCameraService.abortWorkflowRecordings(robotId);
} catch (RuntimeException exception) {
log.warn("清理工作流录像会话失败: robotId={}, error={}", robotId, exception.getMessage());
}
}
private boolean waitForInterruptedNodes(String instId, List<TaskNodeExecuteContext> interruptedNodes) {
long deadlineNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(PAUSE_SETTLE_TIMEOUT_SECONDS);
for (TaskNodeExecuteContext pendingNode : interruptedNodes) {

View File

@ -11,12 +11,15 @@ import lombok.RequiredArgsConstructor;
import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Service;
import java.util.Set;
import java.util.UUID;
@Service
@RequiredArgsConstructor
public class LLMMediaAnalysisOperateService implements LLMOperateService {
private static final Set<String> VIDEO_ANALYSIS_MODES = Set.of("AUTO", "FAST", "ACCURATE");
private final MediaAnalysisClient mediaAnalysisClient;
@Override
@ -48,7 +51,17 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService {
JSONObject options = new JSONObject();
if (!audio) {
options.put("instruction", StringUtils.defaultString(input.getString("instruction")));
options.put("analysisMode", StringUtils.defaultIfBlank(input.getString("analysisMode"), "AUTO"));
String analysisMode = StringUtils.defaultIfBlank(input.getString("analysisMode"), "AUTO")
.trim().toUpperCase();
if (!VIDEO_ANALYSIS_MODES.contains(analysisMode)) {
throw new GlobalException("视频分析模式无效: {}", analysisMode);
}
options.put("analysisMode", analysisMode);
options.put("decisionPolicy", "FAIL_CLOSED");
JSONObject tuning = buildVideoTuning(input.getJSONObject("analysisTuning"));
if (!tuning.isEmpty()) {
options.put("tuning", tuning);
}
}
JSONObject request = new JSONObject();
@ -75,6 +88,31 @@ public class LLMMediaAnalysisOperateService implements LLMOperateService {
return TaskNodeExecuteResult.success(output);
}
private JSONObject buildVideoTuning(JSONObject source) {
JSONObject tuning = new JSONObject();
if (source == null || source.isEmpty()) {
return tuning;
}
Integer maxFrames = source.getInteger("maxFrames");
Integer maxWidth = source.getInteger("maxWidth");
Double confidenceThreshold = source.getDouble("confidenceThreshold");
Boolean fallbackToAccurate = source.getBoolean("fallbackToAccurate");
validateRange("抽帧数量", maxFrames, 2, 16);
validateRange("图片宽度", maxWidth, 640, 1280);
validateRange("通过置信度", confidenceThreshold, 0.5D, 0.95D);
if (maxFrames != null) tuning.put("maxFrames", maxFrames);
if (maxWidth != null) tuning.put("maxWidth", maxWidth);
if (confidenceThreshold != null) tuning.put("confidenceThreshold", confidenceThreshold);
if (fallbackToAccurate != null) tuning.put("fallbackToAccurate", fallbackToAccurate);
return tuning;
}
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);
}
}
private String buildRequestId(TaskNodeExecuteMessage message) {
String instance = StringUtils.defaultIfBlank(message.getInstId(), "single");
String node = StringUtils.defaultIfBlank(message.getNodeId(), "node");

View File

@ -7,9 +7,14 @@ import com.cmvr.test.flow.context.TaskThreadRegistry;
import com.cmvr.test.model.domain.TeTaskInst;
import com.cmvr.test.service.ITeNodeInstService;
import com.cmvr.test.service.ITeTaskInstService;
import com.cmvr.edge.client.service.EdgeCameraService;
import com.cmvr.edge.client.service.EdgeMicrophoneService;
import com.cmvr.edge.client.service.EdgeSystemService;
import org.junit.Test;
import java.lang.reflect.Proxy;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
@ -36,7 +41,8 @@ public class FlowControlServiceTest {
}
};
FlowControlService service = new FlowControlService(
null, holder, nodeInstService, new TaskThreadRegistry(), null, null, taskInstService);
null, holder, nodeInstService, new TaskThreadRegistry(), null,
null, null, null, taskInstService);
service.stop("inst-1");
@ -44,6 +50,54 @@ public class FlowControlServiceTest {
assertTrue(unfinishedNodesStopped.get());
}
@Test
public void clearsLocalMediaSessionsBeforeStoppingRobot() {
List<String> calls = new ArrayList<>();
TaskContext context = new TaskContext();
context.setInstId("inst-1");
context.setRobotId("robot-1");
context.setRobotIds(List.of("robot-1"));
TaskInstHolder holder = new TaskInstHolder(null, null, null, null) {
@Override
public TaskContext getContext(String instId) {
return context;
}
@Override
public void syncStatus(String instId, TaskStatusEnum status) {
context.setStatus(status);
context.setStopped(TaskStatusEnum.STOPPED == status);
}
@Override
public void unRegister(String instId) {
}
};
EdgeMicrophoneService microphoneService = proxy(EdgeMicrophoneService.class, (method, args) -> {
if ("abortWorkflowRecordings".equals(method.getName())) {
assertTrue(context.isStopped());
calls.add("microphone");
}
return defaultValue(method.getReturnType());
});
EdgeCameraService cameraService = proxy(EdgeCameraService.class, (method, args) -> {
if ("abortWorkflowRecordings".equals(method.getName())) calls.add("camera");
return defaultValue(method.getReturnType());
});
EdgeSystemService systemService = proxy(EdgeSystemService.class, (method, args) -> {
if ("stopAll".equals(method.getName())) calls.add("system");
return defaultValue(method.getReturnType());
});
FlowControlService service = new FlowControlService(
null, holder, null, new TaskThreadRegistry(), null,
systemService, microphoneService, cameraService, null);
service.stop("inst-1");
assertEquals(List.of("microphone", "camera", "system"), calls);
assertEquals(TaskStatusEnum.STOPPED, context.getStatus());
}
private ITeTaskInstService taskInstService(String status) {
return (ITeTaskInstService) Proxy.newProxyInstance(
ITeTaskInstService.class.getClassLoader(),
@ -71,4 +125,22 @@ public class FlowControlServiceTest {
throw new UnsupportedOperationException(method.getName());
});
}
@SuppressWarnings("unchecked")
private <T> T proxy(Class<T> type, Invocation invocation) {
return (T) Proxy.newProxyInstance(type.getClassLoader(), new Class<?>[]{type},
(proxy, method, args) -> invocation.invoke(method, args));
}
private Object defaultValue(Class<?> type) {
if (!type.isPrimitive()) return null;
if (type == boolean.class) return false;
if (type == char.class) return '\0';
return 0;
}
@FunctionalInterface
private interface Invocation {
Object invoke(java.lang.reflect.Method method, Object[] args) throws Throwable;
}
}

View File

@ -2,6 +2,7 @@ package com.cmvr.test.flow.runtime.operator.llm;
import com.alibaba.fastjson2.JSONArray;
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;
@ -64,6 +65,11 @@ 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", 8)
.fluentPut("maxWidth", 960)
.fluentPut("confidenceThreshold", 0.8D)
.fluentPut("fallbackToAccurate", false))
.fluentPut("instruction", "检查仪表是否点亮");
service.execute(message(ActionEnum.VIDEO_ANALYZE, input));
@ -72,6 +78,25 @@ public class LLMMediaAnalysisOperateServiceTest {
assertEquals("检查仪表是否点亮",
captured.get().getJSONObject("options").getString("instruction"));
assertEquals("FAST", captured.get().getJSONObject("options").getString("analysisMode"));
assertEquals("FAIL_CLOSED",
captured.get().getJSONObject("options").getString("decisionPolicy"));
JSONObject tuning = captured.get().getJSONObject("options").getJSONObject("tuning");
assertEquals(8, tuning.getIntValue("maxFrames"));
assertEquals(960, tuning.getIntValue("maxWidth"));
assertEquals(0.8D, tuning.getDoubleValue("confidenceThreshold"), 0.0001D);
assertEquals(false, tuning.getBooleanValue("fallbackToAccurate"));
}
@Test(expected = GlobalException.class)
public void rejectsOutOfRangeVideoTuning() {
LLMMediaAnalysisOperateService service = new LLMMediaAnalysisOperateService(request -> new JSONObject());
JSONObject input = new JSONObject()
.fluentPut("profileCode", "aima.power_video.v1")
.fluentPut("videoUrl", "https://files.example/video.mp4")
.fluentPut("analysisMode", "FAST")
.fluentPut("analysisTuning", new JSONObject().fluentPut("maxFrames", 30));
service.execute(message(ActionEnum.VIDEO_ANALYZE, input));
}
private TaskNodeExecuteMessage message(ActionEnum action, JSONObject input) {