Files
Excel-to-SQL_Server/watch.py
Misaka_Company 493621c440 feat: add Uptime Kuma heartbeat monitoring
- Add heartbeat configuration to config.yaml
- Implement HeartbeatMonitor class in watch.py
- Send heartbeat every 40 seconds to Uptime Kuma push endpoint
- Update config.yaml.example with heartbeat settings
- Add heartbeat logging at INFO level

Co-Authored-By: Claude <noreply@anthropic.com>
2026-06-16 15:05:21 +08:00

211 lines
6.7 KiB
Python

"""
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()