| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219 |
- # -*- 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<String> 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<String> 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<String> resolveMeIds(String deviceIds) {\n"
- " List<Long> 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<String> 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<Long> parseDeviceIds(String deviceIds) {\n"
- " if (StringUtils.isBlank(deviceIds)) {\n"
- " return Collections.emptyList();\n"
- " }\n"
- " List<Long> 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)
|