1.4.2 调整不需要重启Docker点击应用即可生效任务

This commit is contained in:
t@123654
2025-06-26 15:59:49 +08:00
parent 60a2125bd5
commit e4fc0e7b3f
12 changed files with 125 additions and 48 deletions
+10 -1
View File
@@ -184,7 +184,16 @@ async def update_message_task(
except Exception as e:
db.rollback()
return error_response(code=500, message=str(e))
@router.put("/job/fresh",summary="重载任务")
async def fresh_message_task(
current_user: dict = Depends(get_current_user)
):
"""
重载任务
"""
from jobs.mps import reload_job
reload_job()
return success_response(message="任务已经重载成功")
@router.delete("/{task_id}",summary="删除消息任务")
async def delete_message_task(
task_id: int,
+15 -1
View File
@@ -105,7 +105,7 @@ class TaskScheduler:
trigger=trigger,
args=args,
kwargs=kwargs,
id=job_id
id=str(job_id)
)
self._jobs[job.id] = job
logger.info(f"Successfully added job {job.id}")
@@ -128,6 +128,20 @@ class TaskScheduler:
return True
return False
def clear_all_jobs(self) -> int:
"""
清除所有任务
:return: 被删除的任务数量
"""
with self._lock:
job_count = len(self._jobs)
if job_count > 0:
self._scheduler.remove_all_jobs()
self._jobs.clear()
logger.info(f"Removed all {job_count} jobs")
return job_count
def start(self) -> None:
"""启动调度器"""
with self._lock:
+11 -3
View File
@@ -1,7 +1,7 @@
import requests
import json
from core.models import Feed
from driver.wx import DoSuccess
from core.db import DB
from core.models.feed import Feed
from .cfg import cfg,wx_cfg
@@ -25,8 +25,9 @@ class WxGather:
def __init__(self,is_add:bool=False):
self.articles=[]
self.is_add=is_add
self._cookies={}
session= requests.Session()
timeout = (5, 10)
timeout = (5, 10)
session.timeout = timeout
self.session=session
self.get_token()
@@ -113,6 +114,9 @@ class WxGather:
print(f"开始")
self.articles=[]
self.get_token()
if self.token=="" or self.token is None:
self.Error("请先扫码登录公众号平台")
return
import time
self.update_mps(mp_id,Feed(
sync_time=int(time.time()),
@@ -121,6 +125,10 @@ class WxGather:
def Item_Over(self,item=None,CallBack=None):
print(f"item end")
_cookies=[{'name': c.name, 'value': c.value, 'domain': c.domain,'expiry':c.expires,'expires':c.expires} for c in self._cookies]
_cookies.append({'name':'token','value':self.token})
if len(_cookies) > 0:
DoSuccess(_cookies)
if CallBack is not None:
CallBack(item)
pass
@@ -130,10 +138,10 @@ class WxGather:
def Over(self,CallBack=None):
if getattr(self, 'articles', None) is not None:
print(f"成功{len(self.articles)}")
if CallBack is not None:
CallBack(self)
def dateformat(self,timestamp:any):
from datetime import datetime, timezone
# UTC时间对象
+1 -1
View File
@@ -4,7 +4,7 @@ import re
import datetime
from datetime import datetime, timezone
# from core.config import cfg
from .cfg import wx_cfg
from .cfg import wx_cfg,cfg
import core.db as db
def dateformat(timestamp:any):
+1 -1
View File
@@ -80,7 +80,7 @@ class MpsApi(WxGather):
msg = resp.json()
self._cookies=resp.cookies
# 流量控制了, 退出
if msg['base_resp']['ret'] == 200013:
super().Error("frequencey control, stop at {}".format(str(begin)))
+4 -2
View File
@@ -78,7 +78,7 @@ class MpsWeb(WxGather):
resp = session.get(url, headers=self.headers, params = params, verify=False)
msg = resp.json()
self._cookies =resp.cookies
# 流量控制了, 退出
if msg['base_resp']['ret'] == 200013:
super().Error("frequencey control, stop at {}".format(str(begin)))
@@ -87,7 +87,9 @@ class MpsWeb(WxGather):
if msg['base_resp']['ret'] == 200003:
super().Error("Invalid Session, stop at {}".format(str(begin)))
break
if msg['base_resp']['ret'] != 0:
super().Error("错误原因:{}:代码:{}".format(msg['base_resp']['err_msg'],msg['base_resp']['ret']))
break
# 如果返回的内容中为空则结束
if 'publish_page' not in msg:
super().Error("all ariticle parsed")
+26
View File
@@ -0,0 +1,26 @@
import time
def expire(cookies:any) :
if not isinstance(cookies, list):
raise TypeError("cookies参数必须是列表类型")
cookie_expiry=None
for cookie in cookies:
if not isinstance(cookie, dict):
continue
if 'name' in cookie and cookie['name'] == 'slave_sid' and 'expiry' in cookie:
try:
expiry_time = float(cookie['expiry'])
remaining_time = expiry_time - time.time()
if remaining_time > 0:
cookie_expiry = {
'expiry_timestamp': expiry_time,
'remaining_seconds': int(remaining_time),
'expiry_time': time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(expiry_time))
}
break
except ValueError:
print(f"slave_sid 的过期时间戳无效: {cookie['expiry']}")
break
return cookie_expiry
-1
View File
@@ -2,7 +2,6 @@ from .token import set_token
def Success(data):
if data != None:
# print("\n登录结果:")
# print(f"Token: {data['token']}")
print(f"有效时间: {data['expiry']['expiry_time']} (剩余秒数: {data['expiry']['remaining_seconds']})")
set_token(data)
else:
+26 -35
View File
@@ -11,6 +11,7 @@ import os
import re
from threading import Thread
from threading import Timer
from .cookies import expire
import json
from core.print import print_error
class Wx:
@@ -189,10 +190,7 @@ class Wx:
wait = WebDriverWait(controller.driver, 120)
wait.until(EC.url_contains(self.WX_HOME))
self.CallBack=CallBack
self.Call_Success()
self.schedule_refresh(interval=refresh_interval)
except NameError as e:
# 修正此处,确保异常处理逻辑正确
@@ -210,6 +208,25 @@ class Wx:
self.Clean()
controller.close()
return self.SESSION
def format_token(self,cookies:any,token=""):
cookies_str=""
for cookie in cookies:
# print(f"{cookie['name']}={cookie['value']}")
cookies_str+=f"{cookie['name']}={cookie['value']}; "
if token=="":
for cookie in cookies:
if 'token' in cookie['name'].lower():
token= cookie['value']
break
# 计算 slave_sid cookie 有效时间
cookie_expiry = expire(cookies)
return{
'cookies': cookies,
'cookies_str': cookies_str,
'token': token,
'wx_login_url': self.wx_login_url,
'expiry': cookie_expiry
}
def Call_Success(self):
print("登录成功!")
self.HasLogin=True
@@ -219,39 +236,9 @@ class Wx:
# 获取当前所有cookie
cookies = self.controller.driver.get_cookies()
# print("\n获取到的Cookie:")
cookies_str=""
for cookie in cookies:
# print(f"{cookie['name']}={cookie['value']}")
cookies_str+=f"{cookie['name']}={cookie['value']}; "
if token:
print(f"\nToken: {token}")
# 计算 slave_sid cookie 有效时间
cookie_expiry = None
for cookie in cookies:
if cookie['name'] == 'slave_sid' and 'expiry' in cookie:
try:
expiry_time = float(cookie['expiry'])
remaining_time = expiry_time - time.time()
if remaining_time > 0:
cookie_expiry = {
'expiry_timestamp': expiry_time,
'remaining_seconds': int(remaining_time),
'expiry_time': time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(expiry_time))
}
break
except ValueError:
print(f"slave_sid 的过期时间戳无效: {cookie['expiry']}")
break
self.SESSION= {
'cookies': cookies,
'cookies_str': cookies_str,
'token': token,
'wx_login_url': self.wx_login_url,
'expiry': cookie_expiry
}
self.SESSION=self.format_token(cookies,token)
self.HasLogin=True
self.Clean()
# print(cookie_expiry)
if self.CallBack is not None:
self.CallBack(self.SESSION)
@@ -286,4 +273,8 @@ class Wx:
except Exception as e:
print(f"设置cookie过期时出错: {str(e)}")
return False
def DoSuccess(cookies:any) -> dict:
data=WX_API.format_token(cookies)
Success(data)
WX_API = Wx()
+8 -2
View File
@@ -45,7 +45,8 @@ def do_job(mps:list[Feed]=None,task:MessageTask=None):
try:
wx.get_Articles(item.faker_id,CallBack=UpdateArticle,Mps_id=item.id,Mps_title=item.mp_name, MaxPage=1,Over_CallBack=Update_Over,interval=interval)
except Exception as e:
print(e)
print_error(e)
# raise
finally:
count=wx.all_count()
all_count+=count
@@ -70,6 +71,11 @@ def get_feeds(task:MessageTask=None):
mps=wx_db.get_all_mps()
return mps
scheduler=TaskScheduler()
def reload_job():
print_success("重载任务")
scheduler.clear_all_jobs()
start_job()
def start_job():
#开启自动同步未同步 文章任务
from jobs.fetch_no_article import start_sync_content
@@ -89,7 +95,7 @@ def start_job():
cron_exp="* * * * *"
# cron_exp="* * * * * *"
pass
job_id=scheduler.add_cron_job(add_job,cron_expr=cron_exp,args=[get_feeds(task),task])
job_id=scheduler.add_cron_job(add_job,cron_expr=cron_exp,args=[get_feeds(task),task],job_id=str(task.id))
print(f"已添加任务: {job_id}")
scheduler.start()
print("启动任务")
+6
View File
@@ -19,6 +19,12 @@ export const createMessageTask = (data: MessageTaskUpdate) => {
export const updateMessageTask = (id: number, data: MessageTaskUpdate) => {
return http.put(`/wx/message_tasks/${id}`, data)
}
export const FreshJobApi = () => {
return http.put(`/wx/message_tasks/job/fresh`)
}
export const FreshJobByIdApi = (id: number, data: MessageTaskUpdate) => {
return http.put(`/wx/message_tasks/job/fresh/${id}`, data)
}
export const deleteMessageTask = (id: number) => {
return http.delete(`/wx/message_tasks/${id}`)
+17 -1
View File
@@ -1,8 +1,9 @@
<script setup lang="ts">
import { ref, onMounted } from 'vue'
import { listMessageTasks, deleteMessageTask } from '@/api/messageTask'
import { listMessageTasks, deleteMessageTask,FreshJobApi,FreshJobByIdApi } from '@/api/messageTask'
import type { MessageTask } from '@/types/messageTask'
import { useRouter } from 'vue-router'
import { Message } from '@arco-design/web-vue'
const parseCronExpression = (exp: string) => {
const parts = exp.split(' ')
@@ -91,6 +92,12 @@ const handlePageChange = (page: number) => {
const handleAdd = () => {
router.push('/message-tasks/add')
}
const FreshJob = () => {
FreshJobApi().then((data) => {
console.log("刷新任务")
Message.success(data.message||"刷新任务成功")
})
}
const handleEdit = (id: number) => {
router.push(`/message-tasks/edit/${id}`)
@@ -120,6 +127,7 @@ onMounted(() => {
<div class="header">
<h2>消息任务列表</h2>
<a-button type="primary" @click="handleAdd">添加消息任务</a-button>
<a-button type="primary" @click="FreshJob">应用</a-button>
</div>
<a-alert type="info" closable>
注意只有添加了任务消息才会定时执行更新任务
@@ -179,6 +187,14 @@ onMounted(() => {
margin-bottom: 20px;
}
.header h2 {
flex: 1;
}
.header .arco-btn {
margin-left: 10px;
}
h2 {
margin: 0;
color: var(--color-text-1);