main_booker.py 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  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 booker import BuiltinBookerGCO, OrderBookerGCO
  12. from utils.cloud_config import load_remote_or_cache, compute_config_hash, fetch_cloud_config, save_local_cache
  13. def get_gco_class_for_booker(cfg: GroupConfig):
  14. if cfg.booker.account_source == "order":
  15. return OrderBookerGCO
  16. return BuiltinBookerGCO
  17. def sync_booker_groups(
  18. groups_conf: list,
  19. active_wrappers: Dict[str, GCOWrapper],
  20. active_gco_classes: Dict[str, type],
  21. app_logger
  22. ):
  23. new_groups_by_id: Dict[str, GroupConfig] = {}
  24. for item in groups_conf:
  25. cfg = GroupConfig.from_json(item)
  26. new_groups_by_id[cfg.identifier] = cfg
  27. current_ids = list(active_wrappers.keys())
  28. for gid in current_ids:
  29. wrapper = active_wrappers[gid]
  30. new_cfg = new_groups_by_id.get(gid)
  31. if not new_cfg or not new_cfg.enable:
  32. app_logger.info(f"Stopping disabled/removed Booker group [{gid}]...")
  33. try:
  34. wrapper.stop()
  35. except Exception as e:
  36. app_logger.error(f"Error stopping group [{gid}]: {e}")
  37. del active_wrappers[gid]
  38. if gid in active_gco_classes:
  39. del active_gco_classes[gid]
  40. for gid, cfg in new_groups_by_id.items():
  41. if not cfg.enable:
  42. continue
  43. target_gco_class = get_gco_class_for_booker(cfg)
  44. if gid in active_wrappers:
  45. if active_gco_classes.get(gid) != target_gco_class:
  46. app_logger.info(f"Mode changed for Booker group [{gid}]. Restarting group with new class...")
  47. try:
  48. active_wrappers[gid].stop()
  49. except Exception as e:
  50. app_logger.error(f"Error stopping group [{gid}] for class switch: {e}")
  51. try:
  52. wrapper = GCOWrapper(
  53. gco_class=target_gco_class,
  54. gco_cfg=cfg,
  55. redis_conf=REDIS_CFG
  56. )
  57. wrapper.load()
  58. wrapper.start()
  59. active_wrappers[gid] = wrapper
  60. active_gco_classes[gid] = target_gco_class
  61. except Exception as e:
  62. app_logger.error(f"Failed to restart group [{gid}]: {e}")
  63. else:
  64. app_logger.info(f"Updating configuration for active Booker group [{gid}]...")
  65. try:
  66. active_wrappers[gid].update_config(cfg)
  67. except Exception as e:
  68. app_logger.error(f"Error updating config for group [{gid}]: {e}")
  69. else:
  70. mode_str = "ORDER (Bound)" if target_gco_class == OrderBookerGCO else "BUILT-IN (Unbound)"
  71. app_logger.info(f"Starting new wrapper for Booker group [{gid}] (Mode: {mode_str})...")
  72. try:
  73. wrapper = GCOWrapper(
  74. gco_class=target_gco_class,
  75. gco_cfg=cfg,
  76. redis_conf=REDIS_CFG
  77. )
  78. wrapper.load()
  79. wrapper.start()
  80. active_wrappers[gid] = wrapper
  81. active_gco_classes[gid] = target_gco_class
  82. except Exception as e:
  83. app_logger.error(f"Failed to start group [{gid}]: {e}")
  84. def main():
  85. parser = argparse.ArgumentParser(description="Booker Runner")
  86. parser.add_argument(
  87. "-c", "--config",
  88. type=str,
  89. required=False,
  90. default=None,
  91. help="Path to local config.json. If specified, disables cloud config and runs with local file."
  92. )
  93. parser.add_argument(
  94. "--poll-interval",
  95. type=float,
  96. required=False,
  97. default=30.0,
  98. help="Interval in seconds to poll cloud config (default: 30.0s)"
  99. )
  100. args = parser.parse_args()
  101. # 逻辑 1:默认开启云端配置,除非命令行指定了 -c/--config 本地文件路径
  102. if args.config is not None:
  103. config_path = args.config
  104. use_cloud = False
  105. else:
  106. config_path = "config/config_booker.json"
  107. use_cloud = True
  108. # 逻辑 2:直接从环境变量获取 node_id,默认 node01
  109. node_id = os.getenv("NODE_ID") or "node01"
  110. config_key = f"booker_config:{node_id}"
  111. poll_interval = args.poll_interval
  112. app_logger = setup_app_logger("Booker")
  113. app_logger.info(f"Booker Logger is ready! (Node ID: '{node_id}', Cloud Mode: {use_cloud})")
  114. groups_conf, current_hash, loaded_from_remote = load_remote_or_cache(
  115. config_key=config_key,
  116. local_cache_path=config_path,
  117. use_cloud=use_cloud,
  118. logger=app_logger.info
  119. )
  120. active_wrappers: Dict[str, GCOWrapper] = {}
  121. active_gco_classes: Dict[str, type] = {}
  122. sync_booker_groups(groups_conf, active_wrappers, active_gco_classes, app_logger)
  123. app_logger.info(
  124. f"Successfully initialized Booker with {len(active_wrappers)} active group(s)."
  125. )
  126. if use_cloud:
  127. app_logger.info(f"Cloud config active for key '{config_key}'. Polling every {poll_interval}s.")
  128. last_poll_time = time.time()
  129. try:
  130. while True:
  131. time.sleep(1)
  132. now = time.time()
  133. if use_cloud and (now - last_poll_time >= poll_interval):
  134. last_poll_time = now
  135. try:
  136. new_conf = fetch_cloud_config(config_key)
  137. if new_conf is not None:
  138. new_hash = compute_config_hash(new_conf)
  139. if new_hash != current_hash:
  140. app_logger.info(
  141. f"Cloud config change detected (hash: {current_hash[:8]} -> {new_hash[:8]}). Auto-reloading..."
  142. )
  143. save_local_cache(config_path, new_conf)
  144. sync_booker_groups(new_conf, active_wrappers, active_gco_classes, app_logger)
  145. current_hash = new_hash
  146. app_logger.info(f"Auto-reload completed. Currently {len(active_wrappers)} group(s) running.")
  147. except Exception as e:
  148. app_logger.warning(f"Error polling cloud config key '{config_key}': {e}. Retrying next cycle.")
  149. except KeyboardInterrupt:
  150. app_logger.info("Shutting down Bookers...")
  151. for gid, wrapper in list(active_wrappers.items()):
  152. try:
  153. wrapper.stop()
  154. except Exception as e:
  155. app_logger.error(f"Error stopping group [{gid}]: {e}")
  156. if __name__ == "__main__":
  157. main()