mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-01 17:00:21 +08:00
refactor: 重构 CleanupScheduler 为 MaintenanceScheduler,新增 Provider 自动签到任务
- 将 cleanup_scheduler.py 重命名为 maintenance_scheduler.py - CleanupScheduler 重命名为 MaintenanceScheduler,保留向后兼容别名 - 新增 Provider 自动签到定时任务(每天凌晨 1:05 执行) - 前端新增定时任务配置区域,支持开关 Provider 自动签到 - 新增 enable_provider_checkin 系统配置项 Close #112
This commit is contained in:
@@ -463,6 +463,33 @@
|
||||
</div>
|
||||
</CardSection>
|
||||
|
||||
<!-- 定时任务 -->
|
||||
<CardSection
|
||||
title="定时任务"
|
||||
description="配置系统后台定时任务"
|
||||
>
|
||||
<div class="flex items-center gap-6">
|
||||
<div class="flex items-center gap-2">
|
||||
<Switch
|
||||
id="enable-provider-checkin"
|
||||
:model-value="systemConfig.enable_provider_checkin"
|
||||
@update:model-value="handleProviderCheckinToggle"
|
||||
/>
|
||||
<div>
|
||||
<Label
|
||||
for="enable-provider-checkin"
|
||||
class="text-sm cursor-pointer"
|
||||
>
|
||||
启用 Provider 自动签到
|
||||
</Label>
|
||||
<p class="text-xs text-muted-foreground">
|
||||
每天凌晨 1:05 执行
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
</CardSection>
|
||||
|
||||
<!-- 系统版本信息 -->
|
||||
<CardSection
|
||||
title="系统信息"
|
||||
@@ -848,6 +875,8 @@ interface SystemConfig {
|
||||
log_retention_days: number
|
||||
cleanup_batch_size: number
|
||||
audit_log_retention_days: number
|
||||
// 定时任务
|
||||
enable_provider_checkin: boolean
|
||||
}
|
||||
|
||||
const basicConfigLoading = ref(false)
|
||||
@@ -900,6 +929,8 @@ const systemConfig = ref<SystemConfig>({
|
||||
log_retention_days: 365,
|
||||
cleanup_batch_size: 1000,
|
||||
audit_log_retention_days: 30,
|
||||
// 定时任务
|
||||
enable_provider_checkin: true,
|
||||
})
|
||||
|
||||
// 原始配置值(用于检测变动)
|
||||
@@ -1002,6 +1033,8 @@ async function loadSystemConfig() {
|
||||
'log_retention_days',
|
||||
'cleanup_batch_size',
|
||||
'audit_log_retention_days',
|
||||
// 定时任务
|
||||
'enable_provider_checkin',
|
||||
]
|
||||
|
||||
for (const key of configs) {
|
||||
@@ -1134,6 +1167,24 @@ async function handleAutoCleanupToggle(enabled: boolean) {
|
||||
}
|
||||
}
|
||||
|
||||
async function handleProviderCheckinToggle(enabled: boolean) {
|
||||
const previousValue = systemConfig.value.enable_provider_checkin
|
||||
systemConfig.value.enable_provider_checkin = enabled
|
||||
try {
|
||||
await adminApi.updateSystemConfig(
|
||||
'enable_provider_checkin',
|
||||
enabled,
|
||||
'是否启用 Provider 自动签到任务'
|
||||
)
|
||||
success(enabled ? '已启用自动签到' : '已禁用自动签到')
|
||||
} catch (err) {
|
||||
error('保存配置失败')
|
||||
log.error('保存自动签到配置失败:', err)
|
||||
// 回滚状态
|
||||
systemConfig.value.enable_provider_checkin = previousValue
|
||||
}
|
||||
}
|
||||
|
||||
async function saveCleanupConfig() {
|
||||
cleanupConfigLoading.value = true
|
||||
try {
|
||||
|
||||
@@ -651,7 +651,7 @@ class AdminTriggerCleanupAdapter(AdminApiAdapter):
|
||||
|
||||
from sqlalchemy import func
|
||||
|
||||
from src.services.system.cleanup_scheduler import get_cleanup_scheduler
|
||||
from src.services.system.maintenance_scheduler import get_maintenance_scheduler
|
||||
|
||||
db = context.db
|
||||
|
||||
@@ -669,8 +669,8 @@ class AdminTriggerCleanupAdapter(AdminApiAdapter):
|
||||
)
|
||||
|
||||
# 触发清理
|
||||
cleanup_scheduler = get_cleanup_scheduler()
|
||||
await cleanup_scheduler._perform_cleanup()
|
||||
maintenance_scheduler = get_maintenance_scheduler()
|
||||
await maintenance_scheduler._perform_cleanup()
|
||||
|
||||
# 获取清理后的统计信息
|
||||
total_after = db.query(Usage).count()
|
||||
|
||||
28
src/main.py
28
src/main.py
@@ -189,13 +189,13 @@ async def lifespan(app: FastAPI):
|
||||
|
||||
# 启动月卡额度重置调度器(仅一个 worker 执行)
|
||||
logger.info("启动月卡额度重置调度器...")
|
||||
from src.services.system.cleanup_scheduler import get_cleanup_scheduler
|
||||
from src.services.system.maintenance_scheduler import get_maintenance_scheduler
|
||||
from src.services.usage.quota_scheduler import get_quota_scheduler
|
||||
from src.services.model.fetch_scheduler import get_model_fetch_scheduler
|
||||
from src.utils.task_coordinator import StartupTaskCoordinator
|
||||
|
||||
quota_scheduler = get_quota_scheduler()
|
||||
cleanup_scheduler = get_cleanup_scheduler()
|
||||
maintenance_scheduler = get_maintenance_scheduler()
|
||||
model_fetch_scheduler = get_model_fetch_scheduler()
|
||||
task_coordinator = StartupTaskCoordinator(redis_client)
|
||||
|
||||
@@ -207,14 +207,14 @@ async def lifespan(app: FastAPI):
|
||||
logger.info("检测到其他 worker 已运行额度调度器,本实例跳过")
|
||||
quota_scheduler = None
|
||||
|
||||
# 启动清理调度器
|
||||
cleanup_scheduler_active = await task_coordinator.acquire("cleanup_scheduler")
|
||||
if cleanup_scheduler_active:
|
||||
logger.info("启动使用记录清理调度器...")
|
||||
await cleanup_scheduler.start()
|
||||
# 启动维护调度器
|
||||
maintenance_scheduler_active = await task_coordinator.acquire("maintenance_scheduler")
|
||||
if maintenance_scheduler_active:
|
||||
logger.info("启动系统维护调度器...")
|
||||
await maintenance_scheduler.start()
|
||||
else:
|
||||
logger.info("检测到其他 worker 已运行清理调度器,本实例跳过")
|
||||
cleanup_scheduler = None
|
||||
logger.info("检测到其他 worker 已运行维护调度器,本实例跳过")
|
||||
maintenance_scheduler = None
|
||||
|
||||
# 启动模型自动获取调度器
|
||||
model_fetch_scheduler_active = await task_coordinator.acquire("model_fetch_scheduler")
|
||||
@@ -248,11 +248,11 @@ async def lifespan(app: FastAPI):
|
||||
await shutdown_batch_committer()
|
||||
logger.info("[OK] 批量提交器已停止,所有待提交数据已保存")
|
||||
|
||||
# 停止清理调度器
|
||||
if cleanup_scheduler:
|
||||
logger.info("停止使用记录清理调度器...")
|
||||
await cleanup_scheduler.stop()
|
||||
await task_coordinator.release("cleanup_scheduler")
|
||||
# 停止维护调度器
|
||||
if maintenance_scheduler:
|
||||
logger.info("停止系统维护调度器...")
|
||||
await maintenance_scheduler.stop()
|
||||
await task_coordinator.release("maintenance_scheduler")
|
||||
|
||||
# 停止月卡额度重置调度器,并释放分布式锁
|
||||
logger.info("停止月卡额度重置调度器...")
|
||||
|
||||
@@ -6,7 +6,11 @@
|
||||
|
||||
from src.services.system.announcement import AnnouncementService
|
||||
from src.services.system.audit import AuditService
|
||||
from src.services.system.cleanup_scheduler import CleanupScheduler
|
||||
from src.services.system.maintenance_scheduler import (
|
||||
CleanupScheduler, # 兼容旧名称
|
||||
MaintenanceScheduler,
|
||||
get_maintenance_scheduler,
|
||||
)
|
||||
from src.services.system.config import SystemConfigService
|
||||
from src.services.system.scheduler import APP_TIMEZONE, TaskScheduler, get_scheduler
|
||||
from src.services.system.sync_stats import SyncStatsService
|
||||
@@ -15,7 +19,9 @@ __all__ = [
|
||||
"SystemConfigService",
|
||||
"AuditService",
|
||||
"AnnouncementService",
|
||||
"CleanupScheduler",
|
||||
"MaintenanceScheduler",
|
||||
"CleanupScheduler", # 兼容旧名称
|
||||
"get_maintenance_scheduler",
|
||||
"SyncStatsService",
|
||||
"TaskScheduler",
|
||||
"get_scheduler",
|
||||
|
||||
@@ -110,6 +110,10 @@ class SystemConfigService:
|
||||
"value": 1000,
|
||||
"description": "每批次清理的记录数,避免单次操作过大影响数据库性能",
|
||||
},
|
||||
"enable_provider_checkin": {
|
||||
"value": True,
|
||||
"description": "是否启用 Provider 自动签到任务,每天凌晨 1:05 执行",
|
||||
},
|
||||
"provider_priority_mode": {
|
||||
"value": "provider",
|
||||
"description": "优先级策略:provider(提供商优先模式) 或 global_key(全局Key优先模式)",
|
||||
|
||||
@@ -1,14 +1,13 @@
|
||||
"""
|
||||
使用记录清理定时任务
|
||||
系统维护定时任务调度器
|
||||
|
||||
分级清理策略:
|
||||
- detail_log_retention_days: 压缩 request_body 和 response_body 到压缩字段
|
||||
- header_retention_days: 清空 request_headers 和 response_headers
|
||||
- log_retention_days: 删除整条记录
|
||||
|
||||
统计聚合任务:
|
||||
- 每天凌晨聚合前一天的统计数据
|
||||
- 更新全局统计汇总
|
||||
包含以下任务:
|
||||
- 统计聚合:每天凌晨聚合前一天的统计数据
|
||||
- Provider 签到:每天凌晨执行所有已配置 Provider 的签到
|
||||
- 使用记录清理:分级清理策略(压缩、清空、删除)
|
||||
- 审计日志清理:定期清理过期的审计日志
|
||||
- 连接池监控:定期检查数据库连接池状态
|
||||
- Pending 状态清理:清理异常的 Pending 状态记录
|
||||
|
||||
使用 APScheduler 进行任务调度,支持时区配置。
|
||||
"""
|
||||
@@ -21,7 +20,8 @@ from sqlalchemy.orm import Session
|
||||
|
||||
from src.core.logger import logger
|
||||
from src.database import create_session
|
||||
from src.models.database import AuditLog, Usage
|
||||
from src.models.database import AuditLog, Provider, Usage
|
||||
from src.services.provider_ops.service import ProviderOpsService
|
||||
from src.services.system.config import SystemConfigService
|
||||
from src.services.system.scheduler import get_scheduler
|
||||
from src.services.system.stats_aggregator import StatsAggregatorService
|
||||
@@ -29,8 +29,8 @@ from src.services.user.apikey import ApiKeyService
|
||||
from src.utils.compression import compress_json
|
||||
|
||||
|
||||
class CleanupScheduler:
|
||||
"""使用记录清理调度器"""
|
||||
class MaintenanceScheduler:
|
||||
"""系统维护任务调度器"""
|
||||
|
||||
def __init__(self):
|
||||
self.running = False
|
||||
@@ -40,11 +40,11 @@ class CleanupScheduler:
|
||||
async def start(self):
|
||||
"""启动调度器"""
|
||||
if self.running:
|
||||
logger.warning("Cleanup scheduler already running")
|
||||
logger.warning("Maintenance scheduler already running")
|
||||
return
|
||||
|
||||
self.running = True
|
||||
logger.info("使用记录清理调度器已启动")
|
||||
logger.info("系统维护调度器已启动")
|
||||
|
||||
scheduler = get_scheduler()
|
||||
|
||||
@@ -100,6 +100,15 @@ class CleanupScheduler:
|
||||
name="审计日志清理",
|
||||
)
|
||||
|
||||
# Provider 签到任务 - 凌晨 1:05 执行
|
||||
scheduler.add_cron_job(
|
||||
self._scheduled_provider_checkin,
|
||||
hour=1,
|
||||
minute=5, # 在统计聚合任务(1:00)之后执行
|
||||
job_id="provider_checkin",
|
||||
name="Provider签到",
|
||||
)
|
||||
|
||||
# 启动时执行一次初始化任务
|
||||
asyncio.create_task(self._run_startup_tasks())
|
||||
|
||||
@@ -129,7 +138,7 @@ class CleanupScheduler:
|
||||
scheduler = get_scheduler()
|
||||
scheduler.stop()
|
||||
|
||||
logger.info("使用记录清理调度器已停止")
|
||||
logger.info("系统维护调度器已停止")
|
||||
|
||||
# ========== 任务函数(APScheduler 直接调用异步函数) ==========
|
||||
|
||||
@@ -158,6 +167,10 @@ class CleanupScheduler:
|
||||
"""审计日志清理任务(定时调用)"""
|
||||
await self._perform_audit_cleanup()
|
||||
|
||||
async def _scheduled_provider_checkin(self):
|
||||
"""Provider 签到任务(定时调用)"""
|
||||
await self._perform_provider_checkin()
|
||||
|
||||
# ========== 实际任务实现 ==========
|
||||
|
||||
async def _perform_stats_aggregation(self, backfill: bool = False):
|
||||
@@ -465,6 +478,89 @@ class CleanupScheduler:
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
async def _perform_provider_checkin(self):
|
||||
"""执行 Provider 签到任务
|
||||
|
||||
遍历所有已配置 provider_ops 的 Provider,触发签到。
|
||||
签到会在余额查询时一起执行(先签到再查询余额)。
|
||||
"""
|
||||
db = create_session()
|
||||
try:
|
||||
# 检查是否启用签到任务
|
||||
if not SystemConfigService.get_config(db, "enable_provider_checkin", True):
|
||||
logger.info("Provider 签到已禁用,跳过签到任务")
|
||||
return
|
||||
|
||||
# 获取所有已配置 provider_ops 的活跃 Provider(只查询需要的字段)
|
||||
providers = (
|
||||
db.query(Provider.id, Provider.config)
|
||||
.filter(Provider.is_active.is_(True))
|
||||
.all()
|
||||
)
|
||||
provider_ids = [
|
||||
p.id
|
||||
for p in providers
|
||||
if p.config and p.config.get("provider_ops")
|
||||
]
|
||||
|
||||
if not provider_ids:
|
||||
logger.info("无已配置的 Provider,跳过签到任务")
|
||||
return
|
||||
|
||||
logger.info(f"开始执行 Provider 签到,共 {len(provider_ids)} 个...")
|
||||
|
||||
# 创建 ProviderOpsService 并执行批量余额查询(会触发签到)
|
||||
service = ProviderOpsService(db)
|
||||
|
||||
# 使用信号量限制并发,避免同时发起过多请求
|
||||
concurrency = 3 # 签到任务并发数
|
||||
semaphore = asyncio.Semaphore(concurrency)
|
||||
|
||||
async def _checkin_provider(provider_id: str) -> tuple[str, bool, str]:
|
||||
"""执行单个 Provider 的签到"""
|
||||
async with semaphore:
|
||||
try:
|
||||
# 触发余额查询(会先执行签到)
|
||||
result = await service.query_balance(provider_id)
|
||||
# 检查签到结果
|
||||
checkin_success = None
|
||||
checkin_message = ""
|
||||
if result.data and hasattr(result.data, "extra") and result.data.extra:
|
||||
checkin_success = result.data.extra.get("checkin_success")
|
||||
checkin_message = result.data.extra.get("checkin_message", "")
|
||||
if checkin_success is True:
|
||||
return provider_id, True, checkin_message
|
||||
elif checkin_success is False:
|
||||
return provider_id, False, checkin_message
|
||||
else:
|
||||
# None 表示未执行签到(可能没配置 Cookie)
|
||||
return provider_id, False, "未执行签到"
|
||||
except Exception as e:
|
||||
logger.warning(f"Provider {provider_id} 签到失败: {e}")
|
||||
return provider_id, False, str(e)
|
||||
|
||||
# 并行执行签到
|
||||
tasks = [_checkin_provider(pid) for pid in provider_ids]
|
||||
results = await asyncio.gather(*tasks)
|
||||
|
||||
# 统计结果
|
||||
success_count = sum(1 for _, success, _ in results if success)
|
||||
logger.info(
|
||||
f"Provider 签到完成: {success_count}/{len(provider_ids)} 成功"
|
||||
)
|
||||
|
||||
# 记录详细结果
|
||||
for provider_id, success, message in results:
|
||||
if success:
|
||||
logger.debug(f" - {provider_id}: 签到成功 - {message}")
|
||||
elif message != "未执行签到":
|
||||
logger.debug(f" - {provider_id}: 签到失败 - {message}")
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"Provider 签到任务执行失败: {e}")
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
async def _perform_cleanup(self):
|
||||
"""执行清理任务"""
|
||||
db = create_session()
|
||||
@@ -823,12 +919,21 @@ class CleanupScheduler:
|
||||
|
||||
|
||||
# 全局单例
|
||||
_cleanup_scheduler = None
|
||||
_maintenance_scheduler = None
|
||||
|
||||
|
||||
def get_cleanup_scheduler() -> CleanupScheduler:
|
||||
"""获取清理调度器单例"""
|
||||
global _cleanup_scheduler
|
||||
if _cleanup_scheduler is None:
|
||||
_cleanup_scheduler = CleanupScheduler()
|
||||
return _cleanup_scheduler
|
||||
def get_maintenance_scheduler() -> MaintenanceScheduler:
|
||||
"""获取维护调度器单例"""
|
||||
global _maintenance_scheduler
|
||||
if _maintenance_scheduler is None:
|
||||
_maintenance_scheduler = MaintenanceScheduler()
|
||||
return _maintenance_scheduler
|
||||
|
||||
|
||||
# 兼容旧名称(deprecated)
|
||||
def get_cleanup_scheduler() -> MaintenanceScheduler:
|
||||
"""获取维护调度器单例(已废弃,请使用 get_maintenance_scheduler)"""
|
||||
return get_maintenance_scheduler()
|
||||
|
||||
|
||||
CleanupScheduler = MaintenanceScheduler # 兼容旧名称
|
||||
Reference in New Issue
Block a user