""" Watch daemon - periodically syncs LAN Excel files and migrates changed tables to SQL Server. """ import argparse import logging import os import time import threading import urllib.request import urllib.error from logger import setup_logger from migrate import load_config, get_connection, create_schema, migrate_table from sync import sync_all_files, determine_tables_to_migrate logger = logging.getLogger('app') class HeartbeatMonitor: """Uptime Kuma heartbeat monitor running in a separate thread.""" def __init__(self, url, interval_seconds): self.url = url self.interval_seconds = interval_seconds self.running = False self.thread = None logger.info(f"心跳监控初始化: 间隔 {interval_seconds} 秒") def _send_heartbeat(self): """Send a heartbeat request to Uptime Kuma.""" try: # Add current timestamp to ping parameter ping = str(int(time.time())) full_url = f"{self.url}{ping}" request = urllib.request.Request(full_url, method='GET') request.add_header('User-Agent', 'Excel-to-SQL-Server-WatchDaemon/1.0') with urllib.request.urlopen(request, timeout=10) as response: if response.status == 200: logger.info(f"心跳发送成功") else: logger.warning(f"心跳返回异常状态码: {response.status}") except urllib.error.URLError as e: logger.error(f"心跳发送失败 (网络错误): {e}") except Exception as e: logger.error(f"心跳发送失败 (未知错误): {e}") def _run_loop(self): """Main heartbeat loop running in the thread.""" logger.info("心跳监控线程已启动") while self.running: self._send_heartbeat() # Wait for the interval or until stopped for _ in range(self.interval_seconds): if not self.running: break time.sleep(1) logger.info("心跳监控线程已停止") def start(self): """Start the heartbeat monitor thread.""" if self.running: logger.warning("心跳监控已在运行") return self.running = True self.thread = threading.Thread(target=self._run_loop, daemon=True) self.thread.start() logger.info("心跳监控已启动") def stop(self): """Stop the heartbeat monitor thread.""" if not self.running: return logger.info("正在停止心跳监控...") self.running = False # Wait for thread to finish if self.thread and self.thread.is_alive(): self.thread.join(timeout=5) logger.info("心跳监控已停止") def run_once(cfg): """Execute one sync + migrate cycle.""" excel_dir = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'Excel') # 1. Sync files from LAN logger.info("开始同步 LAN 文件...") sync_results = sync_all_files(cfg, excel_dir) # 如果全部同步失败,大概率是网络中断 all_failed = sync_results and all(s == 'error' for s, _ in sync_results.values()) if all_failed: logger.warning( f"所有 LAN 源文件不可达 ({len(sync_results)} 个),可能网络中断,跳过本轮同步" ) return # 2. Determine which tables need migration tables_to_migrate = determine_tables_to_migrate(cfg, sync_results) if not tables_to_migrate: logger.info("无文件变更,跳过迁移") return logger.info(f"需要迁移的表: {', '.join(sorted(tables_to_migrate))}") # 3. Open DB connection and migrate each changed table conn = None try: conn = get_connection(cfg) cursor = conn.cursor() create_schema(cursor, cfg['schema']) for table_cfg in cfg['tables']: table_name = table_cfg['table'] if table_name not in tables_to_migrate: continue try: rows = migrate_table(cursor, table_cfg, cfg, use_truncate=True) logger.info(f"迁移完成: {table_name} ({rows:,} 行)") except Exception as e: logger.error(f"迁移失败: {table_name} - {e}", exc_info=True) # Rollback any uncommitted transaction for this table try: conn.rollback() except Exception: pass cursor.close() except Exception as e: logger.error(f"数据库连接失败: {e}", exc_info=True) finally: if conn: try: conn.close() except Exception: pass def main(): parser = argparse.ArgumentParser(description='Excel to SQL Server Watch Daemon') parser.add_argument('--config', default=None, help='配置文件路径 (默认: config.yaml)') parser.add_argument('--once', action='store_true', help='执行一次后退出') parser.add_argument('--interval', type=int, default=None, help='覆盖配置文件中的 check_interval_minutes') args = parser.parse_args() # Load config cfg = load_config(args.config) # Setup logging log_cfg = cfg.get('logging', {}) setup_logger( 'app', log_dir=log_cfg.get('log_dir'), level=log_cfg.get('level', 'INFO'), max_keep_days=log_cfg.get('max_keep_days', 30), ) interval = args.interval or cfg.get('check_interval_minutes', 30) logger.info("=" * 50) logger.info("Watch daemon 启动") logger.info(f" 同步间隔: {interval} 分钟") logger.info("=" * 50) # Initialize heartbeat monitor if configured heartbeat_monitor = None heartbeat_cfg = cfg.get('heartbeat', {}) if heartbeat_cfg.get('enabled', False): url = heartbeat_cfg.get('url', '').strip() interval_seconds = heartbeat_cfg.get('interval_seconds', 40) if url: heartbeat_monitor = HeartbeatMonitor(url, interval_seconds) heartbeat_monitor.start() else: logger.warning("心跳已启用但 URL 未配置,跳过心跳监控") try: while True: try: run_once(cfg) except Exception as e: logger.error(f"运行周期异常: {e}", exc_info=True) if args.once: logger.info("--once 模式,退出") break logger.info(f"下次检查: {interval} 分钟后") time.sleep(interval * 60) finally: # Stop heartbeat monitor on exit if heartbeat_monitor: heartbeat_monitor.stop() if __name__ == '__main__': main()