Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 20 additions & 7 deletions rq_dashboard/web.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
Worker,
requeue_job,
)
from rq.exceptions import NoSuchJobError
from rq.exceptions import DeserializationError, NoSuchJobError
from rq.job import Job
from rq.registry import (
DeferredJobRegistry,
Expand Down Expand Up @@ -271,11 +271,12 @@ def get_queue_registry_jobs_count(queue_name, registry_name, offset, per_page, o
current_queue = ScheduledJobRegistry(queue_name, connection=connection)
elif registry_name == "canceled":
current_queue = CanceledJobRegistry(queue_name, connection=connection)
else:
current_queue = queue
else:
current_queue = queue
total_items = current_queue.count


if order == 'dsc':
end = total_items - offset
start = max(0, end - per_page)
Expand All @@ -287,10 +288,18 @@ def get_queue_registry_jobs_count(queue_name, registry_name, offset, per_page, o
if order == 'dsc':
job_ids.reverse()

current_queue_jobs = [queue.fetch_job(job_id) for job_id in job_ids]
jobs = [serialize_job(job) for job in current_queue_jobs if job]
corrupt_count = 0
current_queue_jobs = []
for job_id in job_ids:
try:
job = queue.fetch_job(job_id)
if job:
current_queue_jobs.append(job)
except (NoSuchJobError, DeserializationError):
corrupt_count += 1
jobs = [serialize_job(job) for job in current_queue_jobs]

return (total_items, jobs)
return (total_items, jobs, corrupt_count)


def escape_format_instance_list(url_list):
Expand Down Expand Up @@ -508,7 +517,7 @@ def list_jobs(instance_number, queue_name, registry_name, per_page, order, page)
per_page = int(per_page)
offset = (current_page - 1) * per_page

total_items, jobs = get_queue_registry_jobs_count(
total_items, jobs, corrupt_count = get_queue_registry_jobs_count(
queue_name, registry_name, offset, per_page, order, current_app.redis_conn
)

Expand Down Expand Up @@ -594,7 +603,11 @@ def list_jobs(instance_number, queue_name, registry_name, per_page, order, page)
)

return dict(
name=queue_name, registry_name=registry_name, jobs=jobs, pagination=pagination
name=queue_name,
registry_name=registry_name,
jobs=jobs,
pagination=pagination,
corrupt_jobs_count=corrupt_count,
)


Expand Down
22 changes: 22 additions & 0 deletions tests/test_basic.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import json
import time
import unittest
from unittest.mock import patch

import redis
from rq import Queue, Worker
Expand Down Expand Up @@ -88,6 +89,27 @@ def test_registry_jobs_list(self):
self.assertIsInstance(data, dict)
self.assertIn('jobs', data)

def test_list_jobs_includes_corrupt_jobs_count(self):
"""List jobs response includes corrupt_jobs_count (skipped missing/corrupt jobs)."""
response = self.client.get('/0/data/jobs/default/queued/8/asc/1.json')
self.assertEqual(response.status_code, HTTP_OK)
data = json.loads(response.data.decode('utf8'))
self.assertIn('corrupt_jobs_count', data)
self.assertIsInstance(data['corrupt_jobs_count'], int)
self.assertGreaterEqual(data['corrupt_jobs_count'], 0)

def test_list_jobs_corrupt_jobs_count_value(self):
"""When get_queue_registry_jobs_count reports skipped jobs, response exposes the count."""
with patch(
'rq_dashboard.web.get_queue_registry_jobs_count',
return_value=(10, [], 3),
):
response = self.client.get('/0/data/jobs/default/queued/8/asc/1.json')
self.assertEqual(response.status_code, HTTP_OK)
data = json.loads(response.data.decode('utf8'))
self.assertEqual(data['corrupt_jobs_count'], 3)
self.assertEqual(data['jobs'], [])

def test_worker_python_version_field(self):
w = Worker(['q'], connection=self.app.redis_conn)
w.register_birth()
Expand Down
Loading