mirror of
https://github.com/insistence/RuoYi-Vue3-FastAPI.git
synced 2026-09-21 20:55:15 +00:00
feat: 新增跨时区支持并完善定时任务调度 (#127)
* feat: 新增跨时区支持 * perf: 优化代码 * fix: 修复测试 * perf: 优化调度器日志打印
This commit is contained in:
@@ -5,7 +5,7 @@ from fastapi.responses import StreamingResponse
|
||||
from pydantic_validation_decorator import ValidateFields
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from common.annotation.cache_annotation import ApiCache, ApiCacheEvict
|
||||
from common.annotation.cache_annotation import ApiCacheEvict
|
||||
from common.annotation.log_annotation import Log
|
||||
from common.annotation.rate_limit_annotation import ApiRateLimit, ApiRateLimitPreset
|
||||
from common.aspect.db_session import DBSessionDependency
|
||||
@@ -15,7 +15,15 @@ from common.constant import ApiGroup, ApiNamespace
|
||||
from common.enums import BusinessType
|
||||
from common.router import APIRouterPro
|
||||
from common.vo import DataResponseModel, PageResponseModel, ResponseBaseModel
|
||||
from module_admin.entity.vo.job_runtime_vo import (
|
||||
JobExecutionModel,
|
||||
JobExecutionQueryModel,
|
||||
JobMutationResult,
|
||||
JobSyncModel,
|
||||
JobSyncQueryModel,
|
||||
)
|
||||
from module_admin.entity.vo.job_vo import (
|
||||
ChangeJobStatusModel,
|
||||
DeleteJobLogModel,
|
||||
DeleteJobModel,
|
||||
EditJobModel,
|
||||
@@ -23,19 +31,45 @@ from module_admin.entity.vo.job_vo import (
|
||||
JobLogPageQueryModel,
|
||||
JobModel,
|
||||
JobPageQueryModel,
|
||||
JobPreviewRequest,
|
||||
JobPreviewResult,
|
||||
JobRunModel,
|
||||
)
|
||||
from module_admin.entity.vo.user_vo import CurrentUserModel
|
||||
from module_admin.service.job_log_service import JobLogService
|
||||
from module_admin.service.job_service import JobService
|
||||
from utils.common_util import bytes2file_response
|
||||
from utils.cron_util import CronUtil
|
||||
from utils.log_util import logger
|
||||
from utils.response_util import ResponseUtil
|
||||
from utils.time_util import TimezoneUtil
|
||||
|
||||
job_controller = APIRouterPro(
|
||||
prefix='/monitor', order_num=13, tags=['系统监控-定时任务'], dependencies=[PreAuthDependency()]
|
||||
)
|
||||
|
||||
|
||||
@job_controller.post(
|
||||
'/job/preview',
|
||||
summary='预览定时任务执行时刻接口',
|
||||
description=(
|
||||
'根据 Quartz Cron 表达式和任务 IANA 时区计算未来执行时刻。'
|
||||
'timeZone 省略时使用系统业务时区,startTime 省略时使用服务端当前时刻;'
|
||||
'count 默认为 5,支持 1~20。返回 UTC 起算时刻及执行时刻列表,无后续匹配时返回空列表。'
|
||||
),
|
||||
response_model=DataResponseModel[JobPreviewResult],
|
||||
dependencies=[UserInterfaceAuthDependency(['monitor:job:add', 'monitor:job:edit'])],
|
||||
)
|
||||
async def preview_system_job(request: Request, preview: JobPreviewRequest) -> Response:
|
||||
start_time = preview.start_time or TimezoneUtil.utc_now()
|
||||
result = JobPreviewResult(
|
||||
timeZone=preview.time_zone,
|
||||
startTime=start_time,
|
||||
nextRunTimes=CronUtil.next_run_times(preview.cron_expression, preview.time_zone, start_time, preview.count),
|
||||
)
|
||||
return ResponseUtil.success(data=result)
|
||||
|
||||
|
||||
@job_controller.get(
|
||||
'/job/list',
|
||||
summary='获取定时任务分页列表接口',
|
||||
@@ -43,7 +77,6 @@ job_controller = APIRouterPro(
|
||||
response_model=PageResponseModel[JobModel],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:list')],
|
||||
)
|
||||
@ApiCache(namespace=ApiNamespace.MONITOR_JOB_LIST)
|
||||
async def get_system_job_list(
|
||||
request: Request,
|
||||
job_page_query: Annotated[JobPageQueryModel, Query()],
|
||||
@@ -60,7 +93,7 @@ async def get_system_job_list(
|
||||
'/job',
|
||||
summary='新增定时任务接口',
|
||||
description='用于新增定时任务',
|
||||
response_model=ResponseBaseModel,
|
||||
response_model=DataResponseModel[JobMutationResult],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:add')],
|
||||
)
|
||||
@ValidateFields(validate_model='add_job')
|
||||
@@ -77,14 +110,14 @@ async def add_system_job(
|
||||
add_job_result = await JobService.add_job_services(query_db, add_job)
|
||||
logger.info(add_job_result.message)
|
||||
|
||||
return ResponseUtil.success(msg=add_job_result.message)
|
||||
return ResponseUtil.success(msg=add_job_result.message, data=add_job_result.result)
|
||||
|
||||
|
||||
@job_controller.put(
|
||||
'/job',
|
||||
summary='编辑定时任务接口',
|
||||
description='用于编辑定时任务',
|
||||
response_model=ResponseBaseModel,
|
||||
response_model=DataResponseModel[JobMutationResult],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:edit')],
|
||||
)
|
||||
@ValidateFields(validate_model='edit_job')
|
||||
@@ -100,21 +133,21 @@ async def edit_system_job(
|
||||
edit_job_result = await JobService.edit_job_services(query_db, edit_job)
|
||||
logger.info(edit_job_result.message)
|
||||
|
||||
return ResponseUtil.success(msg=edit_job_result.message)
|
||||
return ResponseUtil.success(msg=edit_job_result.message, data=edit_job_result.result)
|
||||
|
||||
|
||||
@job_controller.put(
|
||||
'/job/changeStatus',
|
||||
summary='修改定时任务状态接口',
|
||||
description='用于修改定时任务状态',
|
||||
response_model=ResponseBaseModel,
|
||||
response_model=DataResponseModel[JobMutationResult],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:changeStatus')],
|
||||
)
|
||||
@ApiCacheEvict(namespaces=ApiGroup.JOB_MUTATION)
|
||||
@Log(title='定时任务', business_type=BusinessType.UPDATE)
|
||||
async def change_system_job_status(
|
||||
request: Request,
|
||||
change_job: EditJobModel,
|
||||
change_job: ChangeJobStatusModel,
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
current_user: Annotated[CurrentUserModel, CurrentUserDependency()],
|
||||
) -> Response:
|
||||
@@ -127,14 +160,14 @@ async def change_system_job_status(
|
||||
edit_job_result = await JobService.edit_job_services(query_db, edit_job)
|
||||
logger.info(edit_job_result.message)
|
||||
|
||||
return ResponseUtil.success(msg=edit_job_result.message)
|
||||
return ResponseUtil.success(msg=edit_job_result.message, data=edit_job_result.result)
|
||||
|
||||
|
||||
@job_controller.put(
|
||||
'/job/run',
|
||||
summary='执行定时任务接口',
|
||||
description='用于执行指定的定时任务',
|
||||
response_model=ResponseBaseModel,
|
||||
response_model=DataResponseModel[JobExecutionModel],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:changeStatus')],
|
||||
)
|
||||
@ApiRateLimit(namespace=ApiNamespace.MONITOR_JOB_RUN, preset=ApiRateLimitPreset.USER_RESOURCE_EXECUTION)
|
||||
@@ -142,20 +175,104 @@ async def change_system_job_status(
|
||||
@Log(title='定时任务', business_type=BusinessType.UPDATE)
|
||||
async def execute_system_job(
|
||||
request: Request,
|
||||
execute_job: JobModel,
|
||||
execute_job: JobRunModel,
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
current_user: Annotated[CurrentUserModel, CurrentUserDependency()],
|
||||
) -> Response:
|
||||
execute_job_result = await JobService.execute_job_once_services(query_db, execute_job)
|
||||
execute_job_result = await JobService.execute_job_once_services(
|
||||
query_db,
|
||||
execute_job,
|
||||
requested_by=current_user.user.user_name,
|
||||
)
|
||||
logger.info(execute_job_result.message)
|
||||
|
||||
return ResponseUtil.success(msg=execute_job_result.message)
|
||||
return ResponseUtil.success(msg=execute_job_result.message, data=execute_job_result.result)
|
||||
|
||||
|
||||
@job_controller.get(
|
||||
'/job/execution/list',
|
||||
summary='获取定时任务执行记录分页列表接口',
|
||||
description=(
|
||||
'按任务ID、执行ID、执行状态和执行来源(manual 手动、cron 定时)分页查询执行记录,'
|
||||
'返回执行结果或未执行原因、计划/开始/结束时刻及耗时。时间字段统一以 UTC 返回,'
|
||||
'支持查询已删除任务保留的执行记录。'
|
||||
),
|
||||
response_model=PageResponseModel[JobExecutionModel],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:query')],
|
||||
)
|
||||
async def get_system_job_executions(
|
||||
request: Request,
|
||||
query: Annotated[JobExecutionQueryModel, Query()],
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
) -> Response:
|
||||
result = await JobService.execution_list_services(query_db, query)
|
||||
return ResponseUtil.success(model_content=result)
|
||||
|
||||
|
||||
@job_controller.get(
|
||||
'/job/execution/{execution_id}',
|
||||
summary='获取定时任务执行记录详情接口',
|
||||
description=(
|
||||
'根据 32 位小写十六进制执行ID查询单次执行的当前状态、执行结果或未执行原因及 UTC 时间信息。'
|
||||
'任务删除后仍可通过执行ID追踪;执行记录不存在时返回业务错误。'
|
||||
),
|
||||
response_model=DataResponseModel[JobExecutionModel],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:query')],
|
||||
)
|
||||
async def get_system_job_execution(
|
||||
request: Request,
|
||||
execution_id: Annotated[str, Path(min_length=32, max_length=32, pattern='^[0-9a-f]{32}$', description='执行ID')],
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
) -> Response:
|
||||
result = await JobService.execution_detail_services(query_db, execution_id)
|
||||
return ResponseUtil.success(data=result)
|
||||
|
||||
|
||||
@job_controller.get(
|
||||
'/job/sync/list',
|
||||
summary='获取定时任务调度同步状态分页列表接口',
|
||||
description=(
|
||||
'按任务ID和同步状态(pending 待同步、applied 已生效、failed 同步失败)分页查询调度同步记录,'
|
||||
'返回最新配置版本、已应用版本、最近同步错误、删除标记及 UTC 应用/更新时间,'
|
||||
'包含任务删除后保留的同步记录。'
|
||||
),
|
||||
response_model=PageResponseModel[JobSyncModel],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:query')],
|
||||
)
|
||||
async def get_system_job_sync_states(
|
||||
request: Request,
|
||||
query: Annotated[JobSyncQueryModel, Query()],
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
) -> Response:
|
||||
result = await JobService.sync_list_services(query_db, query)
|
||||
return ResponseUtil.success(model_content=result)
|
||||
|
||||
|
||||
@job_controller.post(
|
||||
'/job/sync/{job_id}',
|
||||
summary='重试定时任务调度同步接口',
|
||||
description=(
|
||||
'根据任务ID重新请求应用已保存的最新任务配置或删除操作,返回提交状态、同步状态及任务同步结果。'
|
||||
'以 syncStatus 判断调度是否生效:pending 待同步、applied 已生效、failed 同步失败;'
|
||||
'失败原因见 syncError。同步记录不存在时返回业务错误。'
|
||||
),
|
||||
response_model=DataResponseModel[JobMutationResult],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:edit')],
|
||||
)
|
||||
async def retry_system_job_sync(
|
||||
request: Request,
|
||||
job_id: Annotated[int, Path(gt=0, description='任务ID')],
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
) -> Response:
|
||||
result = await JobService.retry_sync_services(query_db, job_id)
|
||||
return ResponseUtil.success(msg=result.message, data=result.result)
|
||||
|
||||
|
||||
@job_controller.delete(
|
||||
'/job/{job_ids}',
|
||||
summary='删除定时任务接口',
|
||||
description='用于删除定时任务',
|
||||
response_model=ResponseBaseModel,
|
||||
response_model=DataResponseModel[JobMutationResult],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:remove')],
|
||||
)
|
||||
@ApiRateLimit(namespace=ApiNamespace.MONITOR_JOB_DELETE, preset=ApiRateLimitPreset.USER_DESTRUCTIVE_MUTATION)
|
||||
@@ -170,7 +287,7 @@ async def delete_system_job(
|
||||
delete_job_result = await JobService.delete_job_services(query_db, delete_job)
|
||||
logger.info(delete_job_result.message)
|
||||
|
||||
return ResponseUtil.success(msg=delete_job_result.message)
|
||||
return ResponseUtil.success(msg=delete_job_result.message, data=delete_job_result.result)
|
||||
|
||||
|
||||
@job_controller.get(
|
||||
@@ -180,7 +297,6 @@ async def delete_system_job(
|
||||
response_model=DataResponseModel[JobModel],
|
||||
dependencies=[UserInterfaceAuthDependency('monitor:job:query')],
|
||||
)
|
||||
@ApiCache(namespace=ApiNamespace.MONITOR_JOB_DETAIL)
|
||||
async def query_detail_system_job(
|
||||
request: Request,
|
||||
job_id: Annotated[int, Path(description='任务ID')],
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import uuid
|
||||
from datetime import datetime, timedelta
|
||||
from datetime import timedelta
|
||||
from typing import Annotated
|
||||
|
||||
from fastapi import Depends, Request, Response
|
||||
@@ -29,6 +29,7 @@ from module_admin.service.user_service import UserService
|
||||
from utils.jwt_util import JwtUtil
|
||||
from utils.log_util import logger
|
||||
from utils.response_util import ResponseUtil
|
||||
from utils.time_util import TimezoneUtil
|
||||
|
||||
login_controller = APIRouterPro(order_num=1, tags=['登录模块'])
|
||||
|
||||
@@ -85,7 +86,7 @@ async def login(
|
||||
ex=timedelta(minutes=JwtConfig.jwt_redis_expire_minutes),
|
||||
)
|
||||
await UserService.edit_user_services(
|
||||
query_db, EditUserModel(userId=result[0].user_id, loginDate=datetime.now(), type='status')
|
||||
query_db, EditUserModel(userId=result[0].user_id, loginDate=TimezoneUtil.utc_now(), type='status')
|
||||
)
|
||||
logger.info('登录成功')
|
||||
# 判断请求是否来自于api文档,如果是返回指定格式的结果,用于修复api文档认证成功后token显示undefined的bug
|
||||
|
||||
@@ -26,12 +26,6 @@ transport_crypto_controller = APIRouterPro(prefix='/transport/crypto', order_num
|
||||
)
|
||||
@ApiRateLimit(namespace=ApiNamespace.TRANSPORT_CRYPTO_FRONTEND_CONFIG, preset=ApiRateLimitPreset.ANON_PUBLIC_METADATA)
|
||||
async def get_transport_frontend_config(request: Request) -> Response:
|
||||
"""
|
||||
获取当前前端传输层加解密运行配置
|
||||
|
||||
:param request: 当前请求对象
|
||||
:return: 前端传输层加解密运行配置响应
|
||||
"""
|
||||
transport_frontend_config = await TransportCryptoService.get_transport_frontend_config_services()
|
||||
logger.info('获取成功')
|
||||
|
||||
@@ -46,12 +40,6 @@ async def get_transport_frontend_config(request: Request) -> Response:
|
||||
)
|
||||
@ApiRateLimit(namespace=ApiNamespace.TRANSPORT_CRYPTO_PUBLIC_KEY, preset=ApiRateLimitPreset.ANON_PUBLIC_METADATA)
|
||||
async def get_transport_public_key(request: Request) -> Response:
|
||||
"""
|
||||
获取当前传输层加密公钥
|
||||
|
||||
:param request: 当前请求对象
|
||||
:return: 公钥下发响应
|
||||
"""
|
||||
transport_public_key = await TransportCryptoService.get_transport_public_key_services()
|
||||
logger.info('获取成功')
|
||||
|
||||
@@ -66,12 +54,6 @@ async def get_transport_public_key(request: Request) -> Response:
|
||||
dependencies=[PreAuthDependency(), UserInterfaceAuthDependency('monitor:transportCrypto:list')],
|
||||
)
|
||||
async def get_transport_crypto_monitor_info(request: Request) -> Response:
|
||||
"""
|
||||
获取基于Redis聚合的传输层加解密监控信息
|
||||
|
||||
:param request: 当前请求对象
|
||||
:return: 传输层加解密监控信息响应
|
||||
"""
|
||||
transport_crypto_monitor_info = await TransportCryptoService.get_transport_crypto_monitor_info_services(request)
|
||||
logger.info('获取成功')
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import os
|
||||
from datetime import datetime
|
||||
from typing import Annotated, Literal
|
||||
from zoneinfo import available_timezones
|
||||
|
||||
import aiofiles
|
||||
from fastapi import File, Form, Path, Query, Request, Response, UploadFile
|
||||
@@ -41,6 +41,7 @@ from module_admin.entity.vo.user_vo import (
|
||||
UserRoleQueryModel,
|
||||
UserRoleResponseModel,
|
||||
UserRowModel,
|
||||
UserTimezoneModel,
|
||||
)
|
||||
from module_admin.service.dept_service import DeptService
|
||||
from module_admin.service.role_service import RoleService
|
||||
@@ -49,6 +50,7 @@ from utils.common_util import bytes2file_response
|
||||
from utils.log_util import logger
|
||||
from utils.pwd_util import PwdUtil
|
||||
from utils.response_util import ResponseUtil
|
||||
from utils.time_util import TimezoneUtil
|
||||
from utils.upload_util import UploadUtil
|
||||
|
||||
user_controller = APIRouterPro(
|
||||
@@ -222,7 +224,7 @@ async def reset_system_user_pwd(
|
||||
edit_user = EditUserModel(
|
||||
userId=reset_user.user_id,
|
||||
password=PwdUtil.get_password_hash(reset_user.password),
|
||||
pwdUpdateDate=datetime.now(),
|
||||
pwdUpdateDate=TimezoneUtil.utc_now(),
|
||||
updateBy=current_user.user.user_name,
|
||||
type='pwd',
|
||||
)
|
||||
@@ -327,15 +329,16 @@ async def change_system_user_profile_avatar(
|
||||
current_user: Annotated[CurrentUserModel, CurrentUserDependency()],
|
||||
) -> Response:
|
||||
if avatarfile:
|
||||
business_now = TimezoneUtil.to_business_time(TimezoneUtil.utc_now())
|
||||
relative_path = (
|
||||
f'avatar/{datetime.now().strftime("%Y")}/{datetime.now().strftime("%m")}/{datetime.now().strftime("%d")}'
|
||||
f'avatar/{business_now.strftime("%Y")}/{business_now.strftime("%m")}/{business_now.strftime("%d")}'
|
||||
)
|
||||
dir_path = os.path.join(UploadConfig.UPLOAD_PATH, relative_path)
|
||||
try:
|
||||
os.makedirs(dir_path)
|
||||
except FileExistsError:
|
||||
pass
|
||||
avatar_name = f'avatar_{datetime.now().strftime("%Y%m%d%H%M%S")}{UploadConfig.UPLOAD_MACHINE}{UploadUtil.generate_random_number()}.png'
|
||||
avatar_name = f'avatar_{business_now.strftime("%Y%m%d%H%M%S")}{UploadConfig.UPLOAD_MACHINE}{UploadUtil.generate_random_number()}.png'
|
||||
avatar_path = os.path.join(dir_path, avatar_name)
|
||||
async with aiofiles.open(avatar_path, 'wb') as f:
|
||||
await f.write(avatarfile)
|
||||
@@ -385,6 +388,43 @@ async def change_system_user_profile_info(
|
||||
return ResponseUtil.success(msg=edit_user_result.message)
|
||||
|
||||
|
||||
@user_controller.get(
|
||||
'/profile/timezones',
|
||||
summary='获取可选显示时区列表接口',
|
||||
description=(
|
||||
'返回服务端支持的 IANA 时区名称列表,按名称排序并排除 localtime,供当前登录用户选择显示时区。'
|
||||
'如需跟随设备时区,可在修改显示时区接口中提交 timeZone=auto。'
|
||||
),
|
||||
response_model=DataResponseModel[list[str]],
|
||||
)
|
||||
async def get_system_user_timezones(request: Request) -> Response:
|
||||
return ResponseUtil.success(data=sorted(available_timezones() - {'localtime'}))
|
||||
|
||||
|
||||
@user_controller.put(
|
||||
'/profile/timezone',
|
||||
summary='修改当前用户显示时区接口',
|
||||
description=(
|
||||
'保存当前登录账号的显示时区偏好,timeZone 支持 auto(跟随设备)或有效的 IANA 时区名称。'
|
||||
'更新成功后清理用户信息缓存并记录操作日志,返回保存结果;该偏好用于界面时间显示,'
|
||||
'系统业务时区和定时任务时区由各自配置决定。'
|
||||
),
|
||||
response_model=ResponseBaseModel,
|
||||
)
|
||||
@ApiCacheEvict(namespaces=ApiGroup.USER_INFO_MUTATION)
|
||||
@Log(title='时区设置', business_type=BusinessType.UPDATE)
|
||||
async def change_system_user_timezone(
|
||||
request: Request,
|
||||
preference: UserTimezoneModel,
|
||||
query_db: Annotated[AsyncSession, DBSessionDependency()],
|
||||
current_user: Annotated[CurrentUserModel, CurrentUserDependency()],
|
||||
) -> Response:
|
||||
result = await UserService.update_user_timezone_services(
|
||||
query_db, current_user.user.user_id, current_user.user.user_name, preference.time_zone
|
||||
)
|
||||
return ResponseUtil.success(msg=result.message)
|
||||
|
||||
|
||||
@user_controller.put(
|
||||
'/profile/updatePwd',
|
||||
summary='修改用户密码接口',
|
||||
@@ -403,7 +443,7 @@ async def reset_system_user_password(
|
||||
userId=current_user.user.user_id,
|
||||
oldPassword=reset_password.old_password,
|
||||
password=reset_password.new_password,
|
||||
pwdUpdateDate=datetime.now(),
|
||||
pwdUpdateDate=TimezoneUtil.utc_now(),
|
||||
updateBy=current_user.user.user_name,
|
||||
)
|
||||
await UserService.validate_password_services(request.app.state.redis, reset_user.password)
|
||||
|
||||
Reference in New Issue
Block a user