fix_worker_ws_encoding.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. # -*- coding: utf-8 -*-
  2. import pathlib
  3. ROOT = pathlib.Path(__file__).resolve().parents[1]
  4. WORKER_WS_CONSTANTS = (
  5. "package com.fs.common.constant;\n\n"
  6. "/**\n"
  7. " * \u5de5\u4f5c\u673a WebSocket Redis \u901a\u9053\n"
  8. " */\n"
  9. "public class WorkerWsConstants {\n\n"
  10. " /** 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"
  11. " public static final String REDIS_CHANNEL_TASK_RUN = \"worker:ws:task:run\";\n"
  12. "}\n"
  13. )
  14. WORKER_TASK_RUN_WS_MESSAGE = (
  15. "package com.fs.workApp.dto;\n\n"
  16. "import lombok.Data;\n\n"
  17. "import java.util.List;\n\n"
  18. "/**\n"
  19. " * \u5de5\u4f5c\u673a\u4efb\u52a1\u542f\u52a8 WebSocket \u63a8\u9001\uff08Redis Pub/Sub \u6d88\u606f\u4f53\uff09\n"
  20. " */\n"
  21. "@Data\n"
  22. "public class WorkerTaskRunWsMessage {\n\n"
  23. " /** \u5916\u547c\u4efb\u52a1 ID */\n"
  24. " private Long roboticId;\n\n"
  25. " /** \u5916\u547c\u6a21\u5f0f 3/4/5 */\n"
  26. " private Integer callModel;\n\n"
  27. " /** \u76ee\u6807\u8bbe\u5907 IMEI\uff08WS \u8fde\u63a5 meId\uff09 */\n"
  28. " private List<String> meIds;\n"
  29. "}\n"
  30. )
  31. WORKER_TASK_RUN_WS_SUBSCRIBER = (
  32. "package com.fs.worker.websocket.support;\n\n"
  33. "import com.alibaba.fastjson.JSON;\n"
  34. "import com.fs.workApp.dto.WorkerTaskRunWsMessage;\n"
  35. "import com.fs.worker.websocket.service.WorkerWebSocketServer;\n"
  36. "import lombok.extern.slf4j.Slf4j;\n"
  37. "import org.springframework.data.redis.connection.Message;\n"
  38. "import org.springframework.data.redis.connection.MessageListener;\n"
  39. "import org.springframework.stereotype.Component;\n\n"
  40. "import java.nio.charset.StandardCharsets;\n\n"
  41. "/**\n"
  42. " * \u8ba2\u9605 Redis\uff1aPC \u542f\u52a8 APP \u5916\u547c\u4efb\u52a1\u540e\uff0c\u5411\u7ed1\u5b9a\u8bbe\u5907\u63a8\u9001 WebSocket\n"
  43. " */\n"
  44. "@Slf4j\n"
  45. "@Component\n"
  46. "public class WorkerTaskRunWsSubscriber implements MessageListener {\n\n"
  47. " @Override\n"
  48. " public void onMessage(Message message, byte[] pattern) {\n"
  49. " String body = new String(message.getBody(), StandardCharsets.UTF_8);\n"
  50. " try {\n"
  51. " String jsonStr = body;\n"
  52. " if (body.startsWith(\"\\\"\") && body.endsWith(\"\\\"\")) {\n"
  53. " jsonStr = JSON.parseObject(body, String.class);\n"
  54. " }\n"
  55. " WorkerTaskRunWsMessage syncMessage = JSON.parseObject(jsonStr, WorkerTaskRunWsMessage.class);\n"
  56. " if (syncMessage == null || syncMessage.getRoboticId() == null\n"
  57. " || syncMessage.getCallModel() == null || syncMessage.getMeIds() == null) {\n"
  58. " log.warn(\"[WorkerTaskRunWS] \u65e0\u6548\u6d88\u606f body={}\", body);\n"
  59. " return;\n"
  60. " }\n"
  61. " WorkerWebSocketServer.pushTaskRunToDevices(\n"
  62. " syncMessage.getMeIds(),\n"
  63. " syncMessage.getRoboticId(),\n"
  64. " syncMessage.getCallModel());\n"
  65. " } catch (Exception e) {\n"
  66. " log.error(\"[WorkerTaskRunWS] \u5904\u7406\u6d88\u606f\u5931\u8d25 body={}\", body, e);\n"
  67. " }\n"
  68. " }\n"
  69. "}\n"
  70. )
  71. WORKER_WS_REDIS_CONFIG = (
  72. "package com.fs.core.config;\n\n"
  73. "import com.fs.common.constant.WorkerWsConstants;\n"
  74. "import com.fs.worker.websocket.support.WorkerTaskRunWsSubscriber;\n"
  75. "import org.springframework.context.annotation.Bean;\n"
  76. "import org.springframework.context.annotation.Configuration;\n"
  77. "import org.springframework.data.redis.connection.RedisConnectionFactory;\n"
  78. "import org.springframework.data.redis.listener.ChannelTopic;\n"
  79. "import org.springframework.data.redis.listener.RedisMessageListenerContainer;\n\n"
  80. "/**\n"
  81. " * \u5de5\u4f5c\u673a WebSocket Redis \u8ba2\u9605\u914d\u7f6e\n"
  82. " */\n"
  83. "@Configuration\n"
  84. "public class WorkerWsRedisConfig {\n\n"
  85. " @Bean\n"
  86. " public RedisMessageListenerContainer workerWsRedisMessageListenerContainer(\n"
  87. " RedisConnectionFactory redisConnectionFactory,\n"
  88. " WorkerTaskRunWsSubscriber workerTaskRunWsSubscriber) {\n"
  89. " RedisMessageListenerContainer container = new RedisMessageListenerContainer();\n"
  90. " container.setConnectionFactory(redisConnectionFactory);\n"
  91. " container.addMessageListener(workerTaskRunWsSubscriber,\n"
  92. " new ChannelTopic(WorkerWsConstants.REDIS_CHANNEL_TASK_RUN));\n"
  93. " return container;\n"
  94. " }\n"
  95. "}\n"
  96. )
  97. WORKER_TASK_RUN_WS_PUBLISHER = (
  98. "package com.fs.workApp.support;\n\n"
  99. "import com.alibaba.fastjson.JSON;\n"
  100. "import com.fs.common.constant.WorkerWsConstants;\n"
  101. "import com.fs.common.enums.DataSourceType;\n"
  102. "import com.fs.common.utils.StringUtils;\n"
  103. "import com.fs.company.domain.CompanyVoiceRobotic;\n"
  104. "import com.fs.enums.AiCallModeEnum;\n"
  105. "import com.fs.framework.datasource.DynamicDataSourceContextHolder;\n"
  106. "import com.fs.proxy.domain.CompanySmsDevice;\n"
  107. "import com.fs.proxy.service.ICompanySmsPortService;\n"
  108. "import com.fs.workApp.dto.WorkerTaskRunWsMessage;\n"
  109. "import lombok.extern.slf4j.Slf4j;\n"
  110. "import org.springframework.beans.factory.annotation.Autowired;\n"
  111. "import org.springframework.data.redis.core.StringRedisTemplate;\n"
  112. "import org.springframework.stereotype.Component;\n\n"
  113. "import java.util.ArrayList;\n"
  114. "import java.util.Collections;\n"
  115. "import java.util.List;\n\n"
  116. "/**\n"
  117. " * \u4efb\u52a1\u542f\u52a8\u540e\u901a\u8fc7 Redis \u901a\u77e5 fs-worker-app \u5411\u7ed1\u5b9a\u8bbe\u5907\u63a8\u9001 WebSocket \u6d88\u606f\n"
  118. " */\n"
  119. "@Slf4j\n"
  120. "@Component\n"
  121. "public class WorkerTaskRunWsPublisher {\n\n"
  122. " @Autowired(required = false)\n"
  123. " private StringRedisTemplate stringRedisTemplate;\n\n"
  124. " @Autowired\n"
  125. " private ICompanySmsPortService companySmsPortService;\n\n"
  126. " /**\n"
  127. " * 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"
  128. " */\n"
  129. " public void publishIfAppTask(CompanyVoiceRobotic robotic) {\n"
  130. " if (robotic == null || robotic.getId() == null) {\n"
  131. " return;\n"
  132. " }\n"
  133. " if (robotic.getCallModel() == null || !AiCallModeEnum.isAppMode(robotic.getCallModel())) {\n"
  134. " return;\n"
  135. " }\n"
  136. " if (StringUtils.isBlank(robotic.getDeviceIds())) {\n"
  137. " return;\n"
  138. " }\n"
  139. " if (stringRedisTemplate == null) {\n"
  140. " log.warn(\"[WorkerTaskRunWS] StringRedisTemplate \u672a\u6ce8\u5165\uff0c\u8df3\u8fc7\u63a8\u9001 roboticId={}\", robotic.getId());\n"
  141. " return;\n"
  142. " }\n\n"
  143. " List<String> meIds = resolveMeIds(robotic.getDeviceIds());\n"
  144. " if (meIds.isEmpty()) {\n"
  145. " log.warn(\"[WorkerTaskRunWS] \u672a\u89e3\u6790\u5230\u6709\u6548\u8bbe\u5907IMEI roboticId={} deviceIds={}\",\n"
  146. " robotic.getId(), robotic.getDeviceIds());\n"
  147. " return;\n"
  148. " }\n\n"
  149. " WorkerTaskRunWsMessage message = new WorkerTaskRunWsMessage();\n"
  150. " message.setRoboticId(robotic.getId());\n"
  151. " message.setCallModel(robotic.getCallModel());\n"
  152. " message.setMeIds(meIds);\n\n"
  153. " String payload = JSON.toJSONString(message);\n"
  154. " stringRedisTemplate.convertAndSend(WorkerWsConstants.REDIS_CHANNEL_TASK_RUN, payload);\n"
  155. " log.info(\"[WorkerTaskRunWS] \u5df2\u53d1\u5e03 roboticId={} callModel={} meIds={}\",\n"
  156. " robotic.getId(), robotic.getCallModel(), meIds);\n"
  157. " }\n\n"
  158. " private List<String> resolveMeIds(String deviceIds) {\n"
  159. " List<Long> idList = parseDeviceIds(deviceIds);\n"
  160. " if (idList.isEmpty()) {\n"
  161. " return Collections.emptyList();\n"
  162. " }\n\n"
  163. " String previousDataSource = DynamicDataSourceContextHolder.getDataSourceType();\n"
  164. " DynamicDataSourceContextHolder.setDataSourceType(DataSourceType.MASTER.name());\n"
  165. " try {\n"
  166. " List<String> meIds = new ArrayList<>();\n"
  167. " for (Long deviceId : idList) {\n"
  168. " CompanySmsDevice device = companySmsPortService.selectDeviceById(deviceId);\n"
  169. " if (device != null && StringUtils.isNotEmpty(device.getImei())) {\n"
  170. " meIds.add(device.getImei().trim());\n"
  171. " }\n"
  172. " }\n"
  173. " return meIds;\n"
  174. " } finally {\n"
  175. " if (previousDataSource != null) {\n"
  176. " DynamicDataSourceContextHolder.setDataSourceType(previousDataSource);\n"
  177. " } else {\n"
  178. " DynamicDataSourceContextHolder.clearDataSourceType();\n"
  179. " }\n"
  180. " }\n"
  181. " }\n\n"
  182. " private List<Long> parseDeviceIds(String deviceIds) {\n"
  183. " if (StringUtils.isBlank(deviceIds)) {\n"
  184. " return Collections.emptyList();\n"
  185. " }\n"
  186. " List<Long> ids = new ArrayList<>();\n"
  187. " for (String part : deviceIds.split(\",\")) {\n"
  188. " if (StringUtils.isBlank(part)) {\n"
  189. " continue;\n"
  190. " }\n"
  191. " try {\n"
  192. " ids.add(Long.parseLong(part.trim()));\n"
  193. " } catch (NumberFormatException ignored) {\n"
  194. " // \u5ffd\u89c6\u975e\u6cd5\u8bbe\u5907ID\n"
  195. " }\n"
  196. " }\n"
  197. " return ids;\n"
  198. " }\n"
  199. "}\n"
  200. )
  201. FILES = {
  202. ROOT / "fs-common/src/main/java/com/fs/common/constant/WorkerWsConstants.java": WORKER_WS_CONSTANTS,
  203. ROOT / "fs-service/src/main/java/com/fs/workApp/dto/WorkerTaskRunWsMessage.java": WORKER_TASK_RUN_WS_MESSAGE,
  204. ROOT / "fs-service/src/main/java/com/fs/workApp/support/WorkerTaskRunWsPublisher.java": WORKER_TASK_RUN_WS_PUBLISHER,
  205. ROOT / "fs-worker-app/src/main/java/com/fs/worker/websocket/support/WorkerTaskRunWsSubscriber.java": WORKER_TASK_RUN_WS_SUBSCRIBER,
  206. ROOT / "fs-worker-app/src/main/java/com/fs/core/config/WorkerWsRedisConfig.java": WORKER_WS_REDIS_CONFIG,
  207. }
  208. if __name__ == "__main__":
  209. for path, content in FILES.items():
  210. path.write_text(content, encoding="utf-8", newline="\n")
  211. print("fixed:", path)