Skip to content

Commit 8fb4b6e

Browse files
committed
update pydoc
1 parent 19ec6f3 commit 8fb4b6e

File tree

1 file changed

+7
-5
lines changed

1 file changed

+7
-5
lines changed

scheduler/worker/worker.py

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -281,7 +281,7 @@ def handle_job_failure(self, job: JobModel, queue: Queue, exc_string: str = "")
281281
if job.status == JobStatus.FAILED:
282282
self._model.failed_job_count += 1
283283
self._model.completed_jobs += 1
284-
if job.started_at and job.ended_at:
284+
if job.started_at is not None and job.ended_at is not None:
285285
self._model.total_working_time_ms += (job.ended_at - job.started_at).microseconds / 1000.0
286286
self._model.save(connection=queue.connection)
287287

@@ -613,9 +613,11 @@ def monitor_job_execution_process(self, job: JobModel, queue: Queue) -> None:
613613
break
614614
except JobExecutionMonitorTimeoutException:
615615
# job execution process has not exited yet and is still running. Send a heartbeat to keep the worker alive.
616-
working_time = (utcnow() - job.started_at).total_seconds()
617-
self._model.set_current_job_working_time(working_time, self.connection)
618-
616+
if job.started_at is not None:
617+
working_time = (utcnow() - job.started_at).total_seconds()
618+
self._model.set_current_job_working_time(working_time, self.connection)
619+
else:
620+
logger.warning("[Worker {self.name}/{self._pid}]: job.started_at is None, cannot set working time")
619621
# Kill the job from this side if something is really wrong (interpreter lock/etc).
620622
if job.timeout != -1 and self._model.current_job_working_time > (job.timeout + 60):
621623
self._model.heartbeat(self.connection, self.job_monitoring_interval + 60)
@@ -703,7 +705,7 @@ def execute_in_separate_process(self, job: JobModel, queue: Queue) -> None:
703705
random.seed()
704706
self.setup_job_execution_process_signals()
705707
self._is_job_execution_process = True
706-
job = JobModel.get(job.name, self.connection)
708+
job = JobModel.get(job.name, queue.connection)
707709
try:
708710
self.perform_job(job, queue)
709711
except: # noqa

0 commit comments

Comments
 (0)