mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-05 04:58:11 +08:00
Fixes #314 - Convert func to string before saving task in databse so that resubmitting failed task works (#554)
This commit is contained in:
+28
-27
@@ -1,6 +1,7 @@
|
|||||||
# Standard
|
# Standard
|
||||||
import ast
|
import ast
|
||||||
import importlib
|
import inspect
|
||||||
|
import pydoc
|
||||||
import signal
|
import signal
|
||||||
import socket
|
import socket
|
||||||
import traceback
|
import traceback
|
||||||
@@ -416,31 +417,22 @@ def worker(
|
|||||||
f = task["func"]
|
f = task["func"]
|
||||||
# if it's not an instance try to get it from the string
|
# if it's not an instance try to get it from the string
|
||||||
if not callable(task["func"]):
|
if not callable(task["func"]):
|
||||||
try:
|
f = pydoc.locate(f)
|
||||||
module, func = f.rsplit(".", 1)
|
close_old_django_connections()
|
||||||
m = importlib.import_module(module)
|
timer_value = task.pop("timeout", timeout)
|
||||||
f = getattr(m, func)
|
# signal execution
|
||||||
except (ValueError, ImportError, AttributeError) as e:
|
pre_execute.send(sender="django_q", func=f, task=task)
|
||||||
result = (e, False)
|
# execute the payload
|
||||||
if error_reporter:
|
timer.value = timer_value # Busy
|
||||||
error_reporter.report()
|
try:
|
||||||
# We're still going
|
res = f(*task["args"], **task["kwargs"])
|
||||||
if not result:
|
result = (res, True)
|
||||||
close_old_django_connections()
|
except Exception as e:
|
||||||
timer_value = task.pop("timeout", timeout)
|
result = (f"{e} : {traceback.format_exc()}", False)
|
||||||
# signal execution
|
if error_reporter:
|
||||||
pre_execute.send(sender="django_q", func=f, task=task)
|
error_reporter.report()
|
||||||
# execute the payload
|
if task.get("sync", False):
|
||||||
timer.value = timer_value # Busy
|
raise
|
||||||
try:
|
|
||||||
res = f(*task["args"], **task["kwargs"])
|
|
||||||
result = (res, True)
|
|
||||||
except Exception as e:
|
|
||||||
result = (f"{e} : {traceback.format_exc()}", False)
|
|
||||||
if error_reporter:
|
|
||||||
error_reporter.report()
|
|
||||||
if task.get("sync", False):
|
|
||||||
raise
|
|
||||||
with timer.get_lock():
|
with timer.get_lock():
|
||||||
# Process result
|
# Process result
|
||||||
task["result"] = result[0]
|
task["result"] = result[0]
|
||||||
@@ -495,10 +487,19 @@ def save_task(task, broker: Broker):
|
|||||||
broker.acknowledge(task['ack_id'])
|
broker.acknowledge(task['ack_id'])
|
||||||
|
|
||||||
else:
|
else:
|
||||||
|
func = task["func"]
|
||||||
|
# convert func to string
|
||||||
|
if inspect.isfunction(func):
|
||||||
|
func = f"{func.__module__}.{func.__name__}"
|
||||||
|
elif inspect.ismethod(func):
|
||||||
|
func = (
|
||||||
|
f'{func.__self__.__module__}.'
|
||||||
|
f'{func.__self__.__name__}.{func.__name__}'
|
||||||
|
)
|
||||||
Task.objects.create(
|
Task.objects.create(
|
||||||
id=task["id"],
|
id=task["id"],
|
||||||
name=task["name"],
|
name=task["name"],
|
||||||
func=task["func"],
|
func=func,
|
||||||
hook=task.get("hook"),
|
hook=task.get("hook"),
|
||||||
args=task["args"],
|
args=task["args"],
|
||||||
kwargs=task["kwargs"],
|
kwargs=task["kwargs"],
|
||||||
|
|||||||
Reference in New Issue
Block a user