main_standalone.py 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162
  1. import os
  2. import sys
  3. import time
  4. import json
  5. import argparse
  6. from typing import Dict, List, Any
  7. from vs_types import BookerStandaloneConfig
  8. from gco_wrapper import GCOWrapper
  9. from logger_setup import setup_app_logger
  10. from booker_standalone import BookerStandalone
  11. from configure import REDIS_CFG
  12. from utils.cloud_config import load_remote_or_cache, compute_config_hash, fetch_cloud_config, save_local_cache
  13. def sync_standalone_groups(
  14. groups_conf: Any,
  15. active_wrappers: Dict[str, GCOWrapper],
  16. app_logger
  17. ):
  18. """
  19. 根据最新 groups_conf 对比并管理各个 Standalone Booker 实例的生命周期。
  20. 支持 JSON List(包含多个 group 数组)或单个 Dict。
  21. """
  22. if isinstance(groups_conf, dict):
  23. conf_list = [groups_conf]
  24. elif isinstance(groups_conf, list):
  25. conf_list = groups_conf
  26. else:
  27. conf_list = []
  28. new_groups_by_id: Dict[str, BookerStandaloneConfig] = {}
  29. for idx, item in enumerate(conf_list):
  30. cfg = BookerStandaloneConfig.from_json(item)
  31. gid = cfg.identifier or f"standalone_{idx}"
  32. new_groups_by_id[gid] = cfg
  33. # 1. 停止已被移除或被禁用的组
  34. current_ids = list(active_wrappers.keys())
  35. for gid in current_ids:
  36. wrapper = active_wrappers[gid]
  37. new_cfg = new_groups_by_id.get(gid)
  38. if not new_cfg or not new_cfg.enable:
  39. app_logger.info(f"Stopping disabled/removed Standalone group [{gid}]...")
  40. try:
  41. wrapper.stop()
  42. except Exception as e:
  43. app_logger.error(f"Error stopping group [{gid}]: {e}")
  44. del active_wrappers[gid]
  45. # 2. 动态更新现有组,或启动新增组
  46. for gid, cfg in new_groups_by_id.items():
  47. if not cfg.enable:
  48. continue
  49. if gid in active_wrappers:
  50. app_logger.info(f"Updating configuration for active Standalone group [{gid}]...")
  51. try:
  52. active_wrappers[gid].update_config(cfg)
  53. except Exception as e:
  54. app_logger.error(f"Error updating config for Standalone group [{gid}]: {e}")
  55. else:
  56. app_logger.info(f"Starting new wrapper for Standalone group [{gid}]...")
  57. try:
  58. wrapper = GCOWrapper(
  59. gco_class=BookerStandalone,
  60. gco_cfg=cfg,
  61. redis_conf=REDIS_CFG
  62. )
  63. wrapper.load()
  64. wrapper.start()
  65. active_wrappers[gid] = wrapper
  66. except Exception as e:
  67. app_logger.error(f"Failed to start Standalone group [{gid}]: {e}")
  68. def main():
  69. parser = argparse.ArgumentParser(description="Standalone Booker Runner")
  70. parser.add_argument(
  71. "-c", "--config",
  72. type=str,
  73. required=False,
  74. default=None,
  75. help="Path to local config.json. If specified, disables cloud config and runs with local file."
  76. )
  77. parser.add_argument(
  78. "--poll-interval",
  79. type=float,
  80. required=False,
  81. default=30.0,
  82. help="Interval in seconds to poll cloud config (default: 30.0s)"
  83. )
  84. args = parser.parse_args()
  85. # 逻辑 1:默认开启云端配置,除非命令行指定了 -c/--config 本地文件路径
  86. if args.config is not None:
  87. config_path = args.config
  88. use_cloud = False
  89. else:
  90. config_path = "config/config_standalone.json"
  91. use_cloud = True
  92. # 逻辑 2:直接从环境变量获取 node_id,默认 node01
  93. node_id = os.getenv("NODE_ID") or "node01"
  94. config_key = f"standalone_config:{node_id}"
  95. poll_interval = args.poll_interval
  96. app_logger = setup_app_logger("Booker")
  97. app_logger.info(f"Booker Logger is ready! (Node ID: '{node_id}', Cloud Mode: {use_cloud})")
  98. groups_conf, current_hash, loaded_from_remote = load_remote_or_cache(
  99. config_key=config_key,
  100. local_cache_path=config_path,
  101. use_cloud=use_cloud,
  102. logger=app_logger.info
  103. )
  104. active_wrappers: Dict[str, GCOWrapper] = {}
  105. sync_standalone_groups(groups_conf, active_wrappers, app_logger)
  106. app_logger.info(f"Successfully initialized {len(active_wrappers)} Standalone Booker group(s).")
  107. if use_cloud:
  108. app_logger.info(f"Cloud config active for key '{config_key}'. Polling every {poll_interval}s.")
  109. last_poll_time = time.time()
  110. try:
  111. while True:
  112. time.sleep(1)
  113. now = time.time()
  114. if use_cloud and (now - last_poll_time >= poll_interval):
  115. last_poll_time = now
  116. try:
  117. new_conf = fetch_cloud_config(config_key)
  118. if new_conf is not None:
  119. new_hash = compute_config_hash(new_conf)
  120. if new_hash != current_hash:
  121. app_logger.info(
  122. f"Cloud config change detected (hash: {current_hash[:8]} -> {new_hash[:8]}). Auto-reloading..."
  123. )
  124. save_local_cache(config_path, new_conf)
  125. sync_standalone_groups(new_conf, active_wrappers, app_logger)
  126. current_hash = new_hash
  127. app_logger.info(f"Auto-reload completed. Currently {len(active_wrappers)} Standalone group(s) running.")
  128. except Exception as e:
  129. app_logger.warning(f"Error polling cloud config key '{config_key}': {e}. Retrying next cycle.")
  130. except KeyboardInterrupt:
  131. app_logger.info("Shutting down Bookers...")
  132. for gid, wrapper in list(active_wrappers.items()):
  133. try:
  134. wrapper.stop()
  135. except Exception as e:
  136. app_logger.error(f"Error stopping Standalone group [{gid}]: {e}")
  137. if __name__ == "__main__":
  138. main()