Use same FreeGame execution path for scheduler and manual run
This commit is contained in:
@@ -1,7 +1,6 @@
|
|||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
import threading
|
import threading
|
||||||
import traceback
|
import traceback
|
||||||
from datetime import datetime, timedelta
|
|
||||||
|
|
||||||
import requests
|
import requests
|
||||||
from flask import jsonify, render_template
|
from flask import jsonify, render_template
|
||||||
@@ -16,9 +15,6 @@ from .setup import P
|
|||||||
logger = P.logger
|
logger = P.logger
|
||||||
package_name = P.package_name
|
package_name = P.package_name
|
||||||
_fetch_lock = threading.Lock()
|
_fetch_lock = threading.Lock()
|
||||||
_scheduler_lock = threading.Lock()
|
|
||||||
_scheduler_stop_event = threading.Event()
|
|
||||||
_scheduler_thread = None
|
|
||||||
|
|
||||||
SOURCE_LABELS = {
|
SOURCE_LABELS = {
|
||||||
"epic": "Epic",
|
"epic": "Epic",
|
||||||
@@ -57,48 +53,6 @@ def _split_source_payload(results):
|
|||||||
return grouped
|
return grouped
|
||||||
|
|
||||||
|
|
||||||
def _next_run_wait(interval):
|
|
||||||
now = datetime.now()
|
|
||||||
try:
|
|
||||||
from croniter import croniter
|
|
||||||
|
|
||||||
next_run = croniter(interval, now).get_next(datetime)
|
|
||||||
return max(1, int((next_run - now).total_seconds()))
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
|
|
||||||
parts = str(interval or "").split()
|
|
||||||
if len(parts) == 5:
|
|
||||||
minute, hour = parts[0], parts[1]
|
|
||||||
if minute == "*" and hour == "*":
|
|
||||||
return 60
|
|
||||||
if minute.startswith("*/") and hour == "*":
|
|
||||||
try:
|
|
||||||
return max(60, int(minute[2:]) * 60)
|
|
||||||
except Exception:
|
|
||||||
return 60
|
|
||||||
if minute.isdigit() and hour.startswith("*/"):
|
|
||||||
try:
|
|
||||||
target_minute = int(minute)
|
|
||||||
step_hour = int(hour[2:])
|
|
||||||
candidate = now.replace(second=0, microsecond=0)
|
|
||||||
for _ in range(24 * 60):
|
|
||||||
candidate += timedelta(minutes=1)
|
|
||||||
if candidate.minute == target_minute and candidate.hour % step_hour == 0:
|
|
||||||
return max(1, int((candidate - now).total_seconds()))
|
|
||||||
except Exception:
|
|
||||||
return 60
|
|
||||||
return 60
|
|
||||||
|
|
||||||
|
|
||||||
def _scheduler_loop(interval):
|
|
||||||
logger.info("FreeGame internal scheduler loop started: %s", interval)
|
|
||||||
while not _scheduler_stop_event.wait(_next_run_wait(interval)):
|
|
||||||
logger.info("FreeGame internal scheduler tick")
|
|
||||||
Logic.scheduler_function()
|
|
||||||
logger.info("FreeGame internal scheduler loop stopped")
|
|
||||||
|
|
||||||
|
|
||||||
def _discord_send(webhook_url, games):
|
def _discord_send(webhook_url, games):
|
||||||
lines = []
|
lines = []
|
||||||
for game in games[:10]:
|
for game in games[:10]:
|
||||||
@@ -199,7 +153,7 @@ class Logic(PluginModuleBase):
|
|||||||
Logic.scheduler_stop()
|
Logic.scheduler_stop()
|
||||||
return jsonify({"ret": "success"})
|
return jsonify({"ret": "success"})
|
||||||
if sub == "execute_once":
|
if sub == "execute_once":
|
||||||
threading.Thread(target=Logic.scheduler_function, daemon=True).start()
|
Logic.execute_once()
|
||||||
return jsonify({"ret": "success"})
|
return jsonify({"ret": "success"})
|
||||||
if sub == "web_list":
|
if sub == "web_list":
|
||||||
ModelFreeGameItem.ensure_schema()
|
ModelFreeGameItem.ensure_schema()
|
||||||
@@ -214,20 +168,12 @@ class Logic(PluginModuleBase):
|
|||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def scheduler_start():
|
def scheduler_start():
|
||||||
global _scheduler_thread
|
|
||||||
try:
|
try:
|
||||||
interval = ModelSetting.get("auto_interval") or "0 */2 * * *"
|
interval = ModelSetting.get("auto_interval") or "0 */2 * * *"
|
||||||
if F.scheduler.is_include(package_name):
|
if F.scheduler.is_include(package_name):
|
||||||
scheduler.remove_job(package_name)
|
scheduler.remove_job(package_name)
|
||||||
job = Job(package_name, package_name, interval, Logic.scheduler_function, "FreeGame fetch", True)
|
job = Job(package_name, package_name, interval, Logic.execute_once, "FreeGame fetch", True)
|
||||||
scheduler.add_job_instance(job)
|
scheduler.add_job_instance(job)
|
||||||
with _scheduler_lock:
|
|
||||||
if _scheduler_thread is not None and _scheduler_thread.is_alive():
|
|
||||||
_scheduler_stop_event.set()
|
|
||||||
_scheduler_thread.join(timeout=2)
|
|
||||||
_scheduler_stop_event.clear()
|
|
||||||
_scheduler_thread = threading.Thread(target=_scheduler_loop, args=(interval,), daemon=True)
|
|
||||||
_scheduler_thread.start()
|
|
||||||
logger.info("FreeGame scheduler registered: %s", interval)
|
logger.info("FreeGame scheduler registered: %s", interval)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error("Exception:%s", e)
|
logger.error("Exception:%s", e)
|
||||||
@@ -235,22 +181,17 @@ class Logic(PluginModuleBase):
|
|||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def scheduler_stop():
|
def scheduler_stop():
|
||||||
global _scheduler_thread
|
|
||||||
try:
|
try:
|
||||||
if F.scheduler.is_include(package_name):
|
if F.scheduler.is_include(package_name):
|
||||||
scheduler.remove_job(package_name)
|
scheduler.remove_job(package_name)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error("FreeGame framework scheduler stop failed: %s", e)
|
logger.error("FreeGame scheduler stop failed: %s", e)
|
||||||
try:
|
|
||||||
with _scheduler_lock:
|
|
||||||
_scheduler_stop_event.set()
|
|
||||||
if _scheduler_thread is not None and _scheduler_thread.is_alive():
|
|
||||||
_scheduler_thread.join(timeout=2)
|
|
||||||
_scheduler_thread = None
|
|
||||||
except Exception as e:
|
|
||||||
logger.error("Exception:%s", e)
|
|
||||||
logger.error(traceback.format_exc())
|
logger.error(traceback.format_exc())
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def execute_once():
|
||||||
|
threading.Thread(target=Logic.scheduler_function, daemon=True).start()
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def scheduler_function():
|
def scheduler_function():
|
||||||
if not _fetch_lock.acquire(blocking=False):
|
if not _fetch_lock.acquire(blocking=False):
|
||||||
|
|||||||
Reference in New Issue
Block a user