# -*- coding: utf-8 -*- import pathlib ROOT = pathlib.Path(__file__).resolve().parents[1] WORKER_WS_CONSTANTS = ( "package com.fs.common.constant;\n\n" "/**\n" " * \u5de5\u4f5c\u673a WebSocket Redis \u901a\u9053\n" " */\n" "public class WorkerWsConstants {\n\n" " /** PC \u542f\u52a8 APP \u5916\u547c\u4efb\u52a1\u540e\uff0c\u901a\u77e5 fs-worker-app \u5411\u7ed1\u5b9a\u8bbe\u5907\u63a8\u9001 WS \u6d88\u606f */\n" " public static final String REDIS_CHANNEL_TASK_RUN = \"worker:ws:task:run\";\n" "}\n" ) WORKER_TASK_RUN_WS_MESSAGE = ( "package com.fs.workApp.dto;\n\n" "import lombok.Data;\n\n" "import java.util.List;\n\n" "/**\n" " * \u5de5\u4f5c\u673a\u4efb\u52a1\u542f\u52a8 WebSocket \u63a8\u9001\uff08Redis Pub/Sub \u6d88\u606f\u4f53\uff09\n" " */\n" "@Data\n" "public class WorkerTaskRunWsMessage {\n\n" " /** \u5916\u547c\u4efb\u52a1 ID */\n" " private Long roboticId;\n\n" " /** \u5916\u547c\u6a21\u5f0f 3/4/5 */\n" " private Integer callModel;\n\n" " /** \u76ee\u6807\u8bbe\u5907 IMEI\uff08WS \u8fde\u63a5 meId\uff09 */\n" " private List meIds;\n" "}\n" ) WORKER_TASK_RUN_WS_SUBSCRIBER = ( "package com.fs.worker.websocket.support;\n\n" "import com.alibaba.fastjson.JSON;\n" "import com.fs.workApp.dto.WorkerTaskRunWsMessage;\n" "import com.fs.worker.websocket.service.WorkerWebSocketServer;\n" "import lombok.extern.slf4j.Slf4j;\n" "import org.springframework.data.redis.connection.Message;\n" "import org.springframework.data.redis.connection.MessageListener;\n" "import org.springframework.stereotype.Component;\n\n" "import java.nio.charset.StandardCharsets;\n\n" "/**\n" " * \u8ba2\u9605 Redis\uff1aPC \u542f\u52a8 APP \u5916\u547c\u4efb\u52a1\u540e\uff0c\u5411\u7ed1\u5b9a\u8bbe\u5907\u63a8\u9001 WebSocket\n" " */\n" "@Slf4j\n" "@Component\n" "public class WorkerTaskRunWsSubscriber implements MessageListener {\n\n" " @Override\n" " public void onMessage(Message message, byte[] pattern) {\n" " String body = new String(message.getBody(), StandardCharsets.UTF_8);\n" " try {\n" " String jsonStr = body;\n" " if (body.startsWith(\"\\\"\") && body.endsWith(\"\\\"\")) {\n" " jsonStr = JSON.parseObject(body, String.class);\n" " }\n" " WorkerTaskRunWsMessage syncMessage = JSON.parseObject(jsonStr, WorkerTaskRunWsMessage.class);\n" " if (syncMessage == null || syncMessage.getRoboticId() == null\n" " || syncMessage.getCallModel() == null || syncMessage.getMeIds() == null) {\n" " log.warn(\"[WorkerTaskRunWS] \u65e0\u6548\u6d88\u606f body={}\", body);\n" " return;\n" " }\n" " WorkerWebSocketServer.pushTaskRunToDevices(\n" " syncMessage.getMeIds(),\n" " syncMessage.getRoboticId(),\n" " syncMessage.getCallModel());\n" " } catch (Exception e) {\n" " log.error(\"[WorkerTaskRunWS] \u5904\u7406\u6d88\u606f\u5931\u8d25 body={}\", body, e);\n" " }\n" " }\n" "}\n" ) WORKER_WS_REDIS_CONFIG = ( "package com.fs.core.config;\n\n" "import com.fs.common.constant.WorkerWsConstants;\n" "import com.fs.worker.websocket.support.WorkerTaskRunWsSubscriber;\n" "import org.springframework.context.annotation.Bean;\n" "import org.springframework.context.annotation.Configuration;\n" "import org.springframework.data.redis.connection.RedisConnectionFactory;\n" "import org.springframework.data.redis.listener.ChannelTopic;\n" "import org.springframework.data.redis.listener.RedisMessageListenerContainer;\n\n" "/**\n" " * \u5de5\u4f5c\u673a WebSocket Redis \u8ba2\u9605\u914d\u7f6e\n" " */\n" "@Configuration\n" "public class WorkerWsRedisConfig {\n\n" " @Bean\n" " public RedisMessageListenerContainer workerWsRedisMessageListenerContainer(\n" " RedisConnectionFactory redisConnectionFactory,\n" " WorkerTaskRunWsSubscriber workerTaskRunWsSubscriber) {\n" " RedisMessageListenerContainer container = new RedisMessageListenerContainer();\n" " container.setConnectionFactory(redisConnectionFactory);\n" " container.addMessageListener(workerTaskRunWsSubscriber,\n" " new ChannelTopic(WorkerWsConstants.REDIS_CHANNEL_TASK_RUN));\n" " return container;\n" " }\n" "}\n" ) WORKER_TASK_RUN_WS_PUBLISHER = ( "package com.fs.workApp.support;\n\n" "import com.alibaba.fastjson.JSON;\n" "import com.fs.common.constant.WorkerWsConstants;\n" "import com.fs.common.enums.DataSourceType;\n" "import com.fs.common.utils.StringUtils;\n" "import com.fs.company.domain.CompanyVoiceRobotic;\n" "import com.fs.enums.AiCallModeEnum;\n" "import com.fs.framework.datasource.DynamicDataSourceContextHolder;\n" "import com.fs.proxy.domain.CompanySmsDevice;\n" "import com.fs.proxy.service.ICompanySmsPortService;\n" "import com.fs.workApp.dto.WorkerTaskRunWsMessage;\n" "import lombok.extern.slf4j.Slf4j;\n" "import org.springframework.beans.factory.annotation.Autowired;\n" "import org.springframework.data.redis.core.StringRedisTemplate;\n" "import org.springframework.stereotype.Component;\n\n" "import java.util.ArrayList;\n" "import java.util.Collections;\n" "import java.util.List;\n\n" "/**\n" " * \u4efb\u52a1\u542f\u52a8\u540e\u901a\u8fc7 Redis \u901a\u77e5 fs-worker-app \u5411\u7ed1\u5b9a\u8bbe\u5907\u63a8\u9001 WebSocket \u6d88\u606f\n" " */\n" "@Slf4j\n" "@Component\n" "public class WorkerTaskRunWsPublisher {\n\n" " @Autowired(required = false)\n" " private StringRedisTemplate stringRedisTemplate;\n\n" " @Autowired\n" " private ICompanySmsPortService companySmsPortService;\n\n" " /**\n" " * APP \u5916\u547c\u4efb\u52a1\uff08callModel=3/4/5 \u4e14\u5df2\u7ed1\u5b9a\u8bbe\u5907\uff09\u542f\u52a8\u540e\u901a\u77e5\u5de5\u4f5c\u673a\n" " */\n" " public void publishIfAppTask(CompanyVoiceRobotic robotic) {\n" " if (robotic == null || robotic.getId() == null) {\n" " return;\n" " }\n" " if (robotic.getCallModel() == null || !AiCallModeEnum.isAppMode(robotic.getCallModel())) {\n" " return;\n" " }\n" " if (StringUtils.isBlank(robotic.getDeviceIds())) {\n" " return;\n" " }\n" " if (stringRedisTemplate == null) {\n" " log.warn(\"[WorkerTaskRunWS] StringRedisTemplate \u672a\u6ce8\u5165\uff0c\u8df3\u8fc7\u63a8\u9001 roboticId={}\", robotic.getId());\n" " return;\n" " }\n\n" " List meIds = resolveMeIds(robotic.getDeviceIds());\n" " if (meIds.isEmpty()) {\n" " log.warn(\"[WorkerTaskRunWS] \u672a\u89e3\u6790\u5230\u6709\u6548\u8bbe\u5907IMEI roboticId={} deviceIds={}\",\n" " robotic.getId(), robotic.getDeviceIds());\n" " return;\n" " }\n\n" " WorkerTaskRunWsMessage message = new WorkerTaskRunWsMessage();\n" " message.setRoboticId(robotic.getId());\n" " message.setCallModel(robotic.getCallModel());\n" " message.setMeIds(meIds);\n\n" " String payload = JSON.toJSONString(message);\n" " stringRedisTemplate.convertAndSend(WorkerWsConstants.REDIS_CHANNEL_TASK_RUN, payload);\n" " log.info(\"[WorkerTaskRunWS] \u5df2\u53d1\u5e03 roboticId={} callModel={} meIds={}\",\n" " robotic.getId(), robotic.getCallModel(), meIds);\n" " }\n\n" " private List resolveMeIds(String deviceIds) {\n" " List idList = parseDeviceIds(deviceIds);\n" " if (idList.isEmpty()) {\n" " return Collections.emptyList();\n" " }\n\n" " String previousDataSource = DynamicDataSourceContextHolder.getDataSourceType();\n" " DynamicDataSourceContextHolder.setDataSourceType(DataSourceType.MASTER.name());\n" " try {\n" " List meIds = new ArrayList<>();\n" " for (Long deviceId : idList) {\n" " CompanySmsDevice device = companySmsPortService.selectDeviceById(deviceId);\n" " if (device != null && StringUtils.isNotEmpty(device.getImei())) {\n" " meIds.add(device.getImei().trim());\n" " }\n" " }\n" " return meIds;\n" " } finally {\n" " if (previousDataSource != null) {\n" " DynamicDataSourceContextHolder.setDataSourceType(previousDataSource);\n" " } else {\n" " DynamicDataSourceContextHolder.clearDataSourceType();\n" " }\n" " }\n" " }\n\n" " private List parseDeviceIds(String deviceIds) {\n" " if (StringUtils.isBlank(deviceIds)) {\n" " return Collections.emptyList();\n" " }\n" " List ids = new ArrayList<>();\n" " for (String part : deviceIds.split(\",\")) {\n" " if (StringUtils.isBlank(part)) {\n" " continue;\n" " }\n" " try {\n" " ids.add(Long.parseLong(part.trim()));\n" " } catch (NumberFormatException ignored) {\n" " // \u5ffd\u89c6\u975e\u6cd5\u8bbe\u5907ID\n" " }\n" " }\n" " return ids;\n" " }\n" "}\n" ) FILES = { ROOT / "fs-common/src/main/java/com/fs/common/constant/WorkerWsConstants.java": WORKER_WS_CONSTANTS, ROOT / "fs-service/src/main/java/com/fs/workApp/dto/WorkerTaskRunWsMessage.java": WORKER_TASK_RUN_WS_MESSAGE, ROOT / "fs-service/src/main/java/com/fs/workApp/support/WorkerTaskRunWsPublisher.java": WORKER_TASK_RUN_WS_PUBLISHER, ROOT / "fs-worker-app/src/main/java/com/fs/worker/websocket/support/WorkerTaskRunWsSubscriber.java": WORKER_TASK_RUN_WS_SUBSCRIBER, ROOT / "fs-worker-app/src/main/java/com/fs/core/config/WorkerWsRedisConfig.java": WORKER_WS_REDIS_CONFIG, } if __name__ == "__main__": for path, content in FILES.items(): path.write_text(content, encoding="utf-8", newline="\n") print("fixed:", path)