fix_worker_ws_bootstrap_encoding.py 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116
  1. # -*- coding: utf-8 -*-
  2. import pathlib
  3. ROOT = pathlib.Path(r"D:\hdProject\newsaas\ylrz_saas\java\fs-worker-app\src\main\java\com\fs\worker\websocket\support")
  4. BOOTSTRAP = (
  5. "package com.fs.worker.websocket.support;\n\n"
  6. "import com.fs.common.constant.WorkerWsConstants;\n"
  7. "import lombok.extern.slf4j.Slf4j;\n"
  8. "import org.springframework.beans.factory.annotation.Autowired;\n"
  9. "import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;\n"
  10. "import org.springframework.boot.context.event.ApplicationReadyEvent;\n"
  11. "import org.springframework.context.event.EventListener;\n"
  12. "import org.springframework.data.redis.connection.RedisConnectionFactory;\n"
  13. "import org.springframework.data.redis.listener.ChannelTopic;\n"
  14. "import org.springframework.data.redis.listener.RedisMessageListenerContainer;\n"
  15. "import org.springframework.stereotype.Component;\n\n"
  16. "import javax.annotation.PreDestroy;\n\n"
  17. "/**\n"
  18. " * \u5e94\u7528\u5c31\u7eea\u540e\u6ce8\u518c Redis \u8ba2\u9605\uff0c\u907f\u514d\u542f\u52a8\u9636\u6bb5\u8ba2\u9605\u8d85\u65f6\u5bfc\u81f4 fs-worker-app \u65e0\u6cd5\u542f\u52a8\u3002\n"
  19. " */\n"
  20. "@Slf4j\n"
  21. "@Component\n"
  22. "@ConditionalOnBean(RedisConnectionFactory.class)\n"
  23. "public class WorkerWsRedisBootstrap {\n\n"
  24. " private static final long SUBSCRIPTION_WAIT_MS = 30000L;\n\n"
  25. " @Autowired\n"
  26. " private RedisConnectionFactory redisConnectionFactory;\n\n"
  27. " @Autowired\n"
  28. " private WorkerTaskRunWsSubscriber workerTaskRunWsSubscriber;\n\n"
  29. " private RedisMessageListenerContainer container;\n\n"
  30. " @EventListener(ApplicationReadyEvent.class)\n"
  31. " public void onApplicationReady() {\n"
  32. " if (container != null) {\n"
  33. " return;\n"
  34. " }\n"
  35. " try {\n"
  36. " RedisMessageListenerContainer listenerContainer = new RedisMessageListenerContainer();\n"
  37. " listenerContainer.setConnectionFactory(redisConnectionFactory);\n"
  38. " listenerContainer.setMaxSubscriptionRegistrationWaitingTime(SUBSCRIPTION_WAIT_MS);\n"
  39. " listenerContainer.setRecoveryInterval(5000L);\n"
  40. " listenerContainer.addMessageListener(workerTaskRunWsSubscriber,\n"
  41. " new ChannelTopic(WorkerWsConstants.REDIS_CHANNEL_TASK_RUN));\n"
  42. " listenerContainer.afterPropertiesSet();\n"
  43. " listenerContainer.start();\n"
  44. " container = listenerContainer;\n"
  45. " log.info(\"[WorkerTaskRunWS] Redis \u8ba2\u9605\u6210\u529f channel={}\", WorkerWsConstants.REDIS_CHANNEL_TASK_RUN);\n"
  46. " } catch (Exception e) {\n"
  47. " log.warn(\"[WorkerTaskRunWS] Redis \u8ba2\u9605\u5931\u8d25\uff0cPC \u4efb\u52a1\u542f\u52a8 WS \u63a8\u9001\u6682\u4e0d\u53ef\u7528\uff0c\u8bf7\u68c0\u67e5 Redis \u662f\u5426\u542f\u52a8\u4e14\u4e0e fs-saas-company \u540c\u4e00\u5b9e\u4f8b: {}\",\n"
  48. " e.getMessage());\n"
  49. " log.debug(\"[WorkerTaskRunWS] Redis \u8ba2\u9605\u5f02\u5e38\u8be6\u60c5\", e);\n"
  50. " }\n"
  51. " }\n\n"
  52. " @PreDestroy\n"
  53. " public void destroy() {\n"
  54. " if (container != null) {\n"
  55. " try {\n"
  56. " container.stop();\n"
  57. " container.destroy();\n"
  58. " log.info(\"[WorkerTaskRunWS] Redis \u8ba2\u9605\u5df2\u505c\u6b62\");\n"
  59. " } catch (Exception e) {\n"
  60. " log.warn(\"[WorkerTaskRunWS] \u505c\u6b62 Redis \u8ba2\u9605\u5f02\u5e38: {}\", e.getMessage());\n"
  61. " } finally {\n"
  62. " container = null;\n"
  63. " }\n"
  64. " }\n"
  65. " }\n"
  66. "}\n"
  67. )
  68. SUBSCRIBER = (
  69. "package com.fs.worker.websocket.support;\n\n"
  70. "import com.alibaba.fastjson.JSON;\n"
  71. "import com.fs.workApp.dto.WorkerTaskRunWsMessage;\n"
  72. "import com.fs.worker.websocket.service.WorkerWebSocketServer;\n"
  73. "import lombok.extern.slf4j.Slf4j;\n"
  74. "import org.springframework.data.redis.connection.Message;\n"
  75. "import org.springframework.data.redis.connection.MessageListener;\n"
  76. "import org.springframework.stereotype.Component;\n\n"
  77. "import java.nio.charset.StandardCharsets;\n\n"
  78. "/**\n"
  79. " * \u8ba2\u9605 Redis\uff1aPC \u542f\u52a8 APP \u5916\u547c\u4efb\u52a1\u540e\uff0c\u5411\u7ed1\u5b9a\u8bbe\u5907\u63a8\u9001 WebSocket\n"
  80. " */\n"
  81. "@Slf4j\n"
  82. "@Component\n"
  83. "public class WorkerTaskRunWsSubscriber implements MessageListener {\n\n"
  84. " @Override\n"
  85. " public void onMessage(Message message, byte[] pattern) {\n"
  86. " String body = new String(message.getBody(), StandardCharsets.UTF_8);\n"
  87. " try {\n"
  88. " String jsonStr = body;\n"
  89. " if (body.startsWith(\"\\\"\") && body.endsWith(\"\\\"\")) {\n"
  90. " jsonStr = JSON.parseObject(body, String.class);\n"
  91. " }\n"
  92. " WorkerTaskRunWsMessage syncMessage = JSON.parseObject(jsonStr, WorkerTaskRunWsMessage.class);\n"
  93. " if (syncMessage == null || syncMessage.getRoboticId() == null\n"
  94. " || syncMessage.getCallModel() == null || syncMessage.getMeIds() == null) {\n"
  95. " log.warn(\"[WorkerTaskRunWS] \u65e0\u6548\u6d88\u606f body={}\", body);\n"
  96. " return;\n"
  97. " }\n"
  98. " log.info(\"[WorkerTaskRunWS] \u6536\u5230\u6d88\u606f roboticId={} callModel={} meIds={}\",\n"
  99. " syncMessage.getRoboticId(), syncMessage.getCallModel(), syncMessage.getMeIds());\n"
  100. " WorkerWebSocketServer.pushTaskRunToDevices(\n"
  101. " syncMessage.getMeIds(),\n"
  102. " syncMessage.getRoboticId(),\n"
  103. " syncMessage.getCallModel());\n"
  104. " } catch (Exception e) {\n"
  105. " log.error(\"[WorkerTaskRunWS] \u5904\u7406\u6d88\u606f\u5931\u8d25 body={}\", body, e);\n"
  106. " }\n"
  107. " }\n"
  108. "}\n"
  109. )
  110. (ROOT / "WorkerWsRedisBootstrap.java").write_text(BOOTSTRAP, encoding="utf-8", newline="\n")
  111. (ROOT / "WorkerTaskRunWsSubscriber.java").write_text(SUBSCRIBER, encoding="utf-8", newline="\n")
  112. print("done")