main_sentinel.py 5.2 KB

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