CMVR-IOT/cmvr-iot-test/src/main/java/com/cmvr/test/service/FlowiseActionService.java
2026-01-04 17:40:24 +08:00

277 lines
9.6 KiB
Java

package com.cmvr.test.service;
import cn.hutool.http.HttpRequest;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.JSONArray;
import com.alibaba.fastjson2.JSONObject;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.common.utils.http.CallAPIUtil;
import com.cmvr.common.utils.uuid.IdUtils;
import com.cmvr.edge.client.model.EdgeCommonVO;
import com.cmvr.edge.client.service.EdgeBioHeadService;
import com.cmvr.edge.client.service.EdgeHumanoidRobotService;
import com.cmvr.test.enums.FlowiseActionEnum;
import com.cmvr.test.flow.context.TaskContext;
import com.cmvr.test.flow.context.TaskContextManager;
import com.cmvr.test.flow.control.FlowControlService;
import com.cmvr.test.flow.runtime.engine.FlowTaskRuntimeService;
import com.cmvr.test.model.vo.FlowiseActionRequestVO;
import com.cmvr.test.model.vo.FlowiseChatRequestVO;
import com.cmvr.test.model.vo.FlowiseStartRequestVO;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.jetbrains.annotations.NotNull;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.util.HashMap;
import java.util.Map;
@Slf4j
@Service
@RequiredArgsConstructor
public class FlowiseActionService {
private final EdgeBioHeadService edgeBioHeadService;
private final EdgeHumanoidRobotService edgeHumanoidRobotService;
private final TaskContextManager taskContextManager;
private final FlowControlService flowControlService;
private final FlowTaskRuntimeService flowTaskRuntimeService;
@Value("${flowise.tts}")
private String flowiseTts;
@Value("${flowise.query}")
private String flowiseQuery;
@Value("${flowise.start}")
private String flowiseStart;
@Value("${flowise.abort}")
private String flowiseAbort;
@Value("${flowise.api-key}")
private String flowiseApiKey;
@Resource(name = "threadPoolTaskExecutor")
private final ThreadPoolTaskExecutor executor;
private final Object lock = new Object();
private String chatId = null;
public String start(FlowiseStartRequestVO request) {
log.info("flowise start: {}", System.currentTimeMillis());
abort();
synchronized (lock) {
if (chatId != null) {
throw new GlobalException("Flow is already running: " + chatId);
}
chatId = IdUtils.fastUUID();
JSONObject body = new JSONObject();
body.put("question", request.getQuestion());
body.put("chatId", chatId);
body.put("streaming", true);
executor.submit(() -> {
try {
HttpRequest.post(flowiseStart)
.header("Content-Type", "application/json")
.header("Authorization", "Bearer " + flowiseApiKey)
.body(body.toJSONString())
.timeout(0)
.execute()
.body();
} catch (Exception e) {
throw new GlobalException("flowise 流程执行失败: ", e.getMessage());
}
});
return chatId;
}
}
public String abort() {
CallAPIUtil.doPostJson(flowiseTts + "stop", null, null);
Map<String, TaskContext> instContextMap = taskContextManager.getInstContextMap();
instContextMap.forEach((key, value) -> flowControlService.stop(key));
synchronized (lock) {
if (chatId == null) {
return "No flow running.";
}
String id = chatId;
String result = HttpRequest.put(flowiseAbort + id)
.header("Authorization", "Bearer " + flowiseApiKey)
.execute()
.body();
chatId = null;
return result;
}
}
private boolean acOn = false;
private boolean musicOn = false;
public String command(FlowiseActionRequestVO request) {
log.info("command start: {}", System.currentTimeMillis());
String action = request.getAction();
FlowiseActionEnum actionEnum = FlowiseActionEnum.fromAction(action);
String terminalId = "4ed1246c465b97975f96c9ef8371a3bd";
String bioId = "bio_head";
switch (actionEnum) {
// ==== 表情 ====
case SMILE:
log.info("高兴");
return edgeBioHeadService.specialExpression(terminalId, bioId, 1);
case SAD:
log.info("悲伤");
return edgeBioHeadService.specialExpression(terminalId, bioId, 4);
case SURPRISED:
log.info("惊讶");
return edgeBioHeadService.specialExpression(terminalId, bioId, 2);
case TIRED:
log.info("累了");
return edgeBioHeadService.specialExpression(terminalId, bioId, 5);
// ==== 动作 ====
case AC_ON:
if (!acOn) {
log.info("空调打开了");
HttpUtil.post("http://192.168.0.10:8000/action/air", "");
acOn = true;
return "空调打开了";
} else {
// 已经是开着的
speakStatus("空调是开着的呢");
return "空调已经打开";
}
case AC_OFF:
if (acOn) {
log.info("空调关闭了");
HttpUtil.post("http://192.168.0.10:8000/action/air", "");
acOn = false;
return "空调关闭了";
} else {
speakStatus("空调已经是关着的呢");
return "空调已经关闭";
}
case MUSIC_ON:
if (!musicOn) {
log.info("音乐打开了");
HttpUtil.post("http://192.168.0.10:8000/action/music", "");
musicOn = true;
return "音乐打开了";
} else {
speakStatus("音乐是开着的呢");
return "音乐已经打开";
}
case MUSIC_OFF:
if (musicOn) {
log.info("音乐关闭了");
HttpUtil.post("http://192.168.0.10:8000/action/music", "");
musicOn = false;
return "音乐关闭了";
} else {
speakStatus("音乐已经是关着的呢");
return "音乐已经关闭";
}
case DRIVE_MODE_ECONOMY:
HttpUtil.post("http://192.168.0.10:8000/action/drive_economy", "");
log.info("驾驶模式已切换为节能模式");
return "驾驶模式已切换为节能模式";
case DRIVE_MODE_COMFORT:
HttpUtil.post("http://192.168.0.10:8000/action/drive_comfort", "");
log.info("驾驶模式已切换为舒适模式");
return "驾驶模式已切换为舒适模式";
case HELLO:
HttpUtil.post("http://192.168.0.10:8000/action/hello", "");
log.info("打招呼");
return "打招呼";
default:
throw new UnsupportedOperationException("未实现的 Flowise Action: " + action);
}
}
public String chat(FlowiseChatRequestVO requestVO) {
log.info("chat start: {}", System.currentTimeMillis());
Map<String, String> body = new HashMap<>();
body.put("text", requestVO.getText());
body.put("voice", "x4_yezi");
CallAPIUtil.doPostJson(flowiseTts + "play", null, body);
return "ok";
}
public JSONObject query(String chatId) {
String body = HttpRequest.get(flowiseQuery + chatId)
.header("Authorization", "Bearer " + flowiseApiKey)
.execute()
.body();
JSONObject root = JSON.parseObject(body);
String execStr = root.getString("executionData");
JSONArray execJson = JSON.parseArray(execStr);
root.put("executionData", execJson);
return root;
}
public String actionStart() {
EdgeCommonVO speakVO = getEdgeCommonVO();
try {
// 张嘴
String speakStart = edgeBioHeadService.speakStart(speakVO);
log.info("张嘴动作触发完成: {}", speakStart);
} catch (Exception e) {
throw new GlobalException("张嘴执行异常", e);
}
return "ok";
}
@NotNull
private static EdgeCommonVO getEdgeCommonVO() {
String terminalId = "4ed1246c465b97975f96c9ef8371a3bd";
String bioId = "bio_head";
EdgeCommonVO speakVO = new EdgeCommonVO();
speakVO.setTerminalId(terminalId);
speakVO.setDeviceId(bioId);
return speakVO;
}
public String actionStop() {
EdgeCommonVO speakVO = getEdgeCommonVO();
// 停止张嘴
try {
edgeBioHeadService.speakStop(speakVO);
log.info("停止张嘴成功");
} catch (Exception e) {
throw new GlobalException("停止张嘴失败", e);
}
return "ok";
}
private void speakStatus(String text) {
Map<String, String> body = new HashMap<>();
body.put("text", text);
body.put("voice", "x4_yezi");
CallAPIUtil.doPostJson(flowiseTts + "play", null, body);
}
}