More explicit log messages in exception handling (#59)

This commit is contained in:
msabatier
2023-01-11 00:36:35 +01:00
committed by GitHub
parent 31e82ad028
commit 8cd1028391
+15 -10
View File
@@ -366,7 +366,7 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
while True: while True:
try: try:
task_set = broker.dequeue() task_set = broker.dequeue()
except Exception as e: except Exception:
logger.exception("Failed to pull task from broker") logger.exception("Failed to pull task from broker")
# broker probably crashed. Let the sentinel handle it. # broker probably crashed. Let the sentinel handle it.
sleep(10) sleep(10)
@@ -377,7 +377,7 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
# unpack the task # unpack the task
try: try:
task = SignedPackage.loads(task[1]) task = SignedPackage.loads(task[1])
except (TypeError, BadSignature) as e: except (TypeError, BadSignature):
logger.exception("Failed to push task to queue") logger.exception("Failed to push task to queue")
broker.fail(ack_id) broker.fail(ack_id)
continue continue
@@ -473,7 +473,8 @@ 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(f):
# locate() returns None if f cannot be loaded
f = pydoc.locate(f) f = pydoc.locate(f)
close_old_django_connections() close_old_django_connections()
timer_value = task.pop("timeout", timeout) timer_value = task.pop("timeout", timeout)
@@ -482,6 +483,9 @@ def worker(
# execute the payload # execute the payload
timer.value = timer_value # Busy timer.value = timer_value # Busy
try: try:
if f is None:
# raise a meaningfull error if task["func"] is not a valid function
raise ValueError(f"Function {task['func']} is not defined")
res = f(*task["args"], **task["kwargs"]) res = f(*task["args"], **task["kwargs"])
result = (res, True) result = (res, True)
except Exception as e: except Exception as e:
@@ -584,8 +588,8 @@ def save_task(task, broker: Broker):
success=task["success"], success=task["success"],
attempt_count=1, attempt_count=1,
) )
except Exception as e: except Exception:
logger.error(e) logger.exception("Could not save task result")
def save_cached(task, broker: Broker): def save_cached(task, broker: Broker):
@@ -636,8 +640,8 @@ def save_cached(task, broker: Broker):
) )
# save the task # save the task
broker.cache.set(task_key, SignedPackage.dumps(task), timeout) broker.cache.set(task_key, SignedPackage.dumps(task), timeout)
except Exception as e: except Exception:
logger.error(e) logger.exception("Could not save task result")
def scheduler(broker: Broker = None): def scheduler(broker: Broker = None):
@@ -725,11 +729,12 @@ def scheduler(broker: Broker = None):
else: else:
logger.info( logger.info(
_( _(
"%(process_name)s created a task from schedule " "%(process_name)s created task %(task_name)s from schedule "
"[%(schedule)s]" "[%(schedule)s]"
) )
% { % {
"process_name": current_process().name, "process_name": current_process().name,
"task_name": humanize(s.task),
"schedule": s.name or s.id, "schedule": s.name or s.id,
} }
) )
@@ -742,8 +747,8 @@ def scheduler(broker: Broker = None):
s.repeats = 0 s.repeats = 0
# save the schedule # save the schedule
s.save() s.save()
except Exception as e: except Exception:
logger.error(e) logger.exception("Could not create task from schedule")
def close_old_django_connections(): def close_old_django_connections():