Files
domainCheck/app/core/task_scheduler.py
2026-04-14 22:53:52 +08:00

162 lines
4.8 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: UTF-8 -*-
'''
@Project :domainScanDemo
@File :task_scheduler.py
@IDE :PyCharm
@Author :梦伴
@Date :2026/4/8 23:53
@explain : 任务调度器
'''
import time
import threading
from loguru import logger
from app.utils.database import Database
from app.core.detect_engine import DetectEngine
class TaskScheduler:
"""
任务调度器
"""
def __init__(self):
"""
初始化任务调度器
"""
self.db = Database()
self.detect_engine = DetectEngine()
self.running = False
self.threads = []
self.max_threads = 10
def start(self):
"""
启动任务调度器
"""
if self.running:
logger.info("任务调度器已经在运行中")
return
self.running = True
logger.info("启动任务调度器")
# 启动多个线程处理任务
for i in range(self.max_threads):
thread = threading.Thread(target=self._process_tasks, daemon=True)
thread.start()
self.threads.append(thread)
logger.info(f"启动任务处理线程 {i+1}")
def stop(self):
"""
停止任务调度器
"""
self.running = False
logger.info("停止任务调度器")
# 等待线程结束
for thread in self.threads:
thread.join(timeout=5)
self.threads.clear()
logger.info("任务调度器已停止")
def _process_tasks(self):
"""
处理任务
"""
while self.running:
try:
# 获取待执行的任务
task = self.db.get_pending_task()
if not task:
# 没有任务,休眠一段时间
time.sleep(1)
continue
task_id = task['id']
domain_id = task['domain_id']
logger.info(f"处理任务: {task_id}, 域名ID: {domain_id}")
# 执行任务
success = self.detect_engine.process_task(task_id)
if success:
logger.info(f"任务处理成功: {task_id}")
else:
logger.warning(f"任务处理失败: {task_id}")
# 短暂休眠,避免过于频繁的数据库操作
time.sleep(0.1)
except Exception as e:
logger.error(f"处理任务出错: {e}")
# 休眠一段时间,避免出错后无限循环
time.sleep(5)
def add_task(self, domain_id, task_type=1, priority=0):
"""
添加任务
:param domain_id: 域名ID
:param task_type: 任务类型1-基础检测2-深度检测
:param priority: 优先级0-低1-中2-高
:return: int - 任务ID
"""
try:
task_id = self.db.create_detect_task(domain_id, task_type, priority)
logger.info(f"添加任务成功: {task_id}, 域名ID: {domain_id}")
return task_id
except Exception as e:
logger.error(f"添加任务失败: {e}")
return None
def get_task_stats(self):
"""
获取任务统计信息
:return: dict - 任务统计信息
"""
try:
return self.db.get_task_statistics()
except Exception as e:
logger.error(f"获取任务统计信息出错: {e}")
return {}
def retry_failed_tasks(self):
"""
重试失败的任务
:return: int - 重试的任务数量
"""
try:
tasks = self.db.get_failed_tasks()
retry_count = 0
for task in tasks:
task_id = task['id']
self.db.update_task_status(task_id, 0) # 0 表示待执行
self.db.update_task_retry_count(task_id, 0) # 重置重试次数
retry_count += 1
logger.info(f"重试 {retry_count} 个失败的任务")
return retry_count
except Exception as e:
logger.error(f"重试失败任务出错: {e}")
return 0
def clear_completed_tasks(self, days=7):
"""
清理已完成的任务
:param days: 保留天数
:return: int - 清理的任务数量
"""
try:
count = self.db.clear_completed_tasks(days)
logger.info(f"清理 {count} 个已完成的任务")
return count
except Exception as e:
logger.error(f"清理已完成任务出错: {e}")
return 0