Skip to content

Commit 4c38cbc

Browse files
committed
yaas
1 parent 4a7f085 commit 4c38cbc

2 files changed

Lines changed: 60 additions & 1 deletion

File tree

app/api/urls.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
OpsQueuePeekView,
3535
OpsQueuePurgeView,
3636
OpsQueueEvictPurgeView,
37+
OpsConnectionCloseView,
3738
OpsTaskForceCloseView,
3839
)
3940
from .views.analytics import (
@@ -176,6 +177,7 @@
176177
re_path(r'^ops/queues/(?P<name>[^/]+)/peek/?$', OpsQueuePeekView.as_view(), name='ops-queue-peek'),
177178
re_path(r'^ops/queues/(?P<name>[^/]+)/purge/?$', OpsQueuePurgeView.as_view(), name='ops-queue-purge'),
178179
re_path(r'^ops/queues/?$', OpsQueuesView.as_view(), name='ops-queues'),
180+
re_path(r'^ops/connections/close/?$', OpsConnectionCloseView.as_view(), name='ops-connection-close'),
179181
re_path(r'^ops/tasks/(?P<task_id>[0-9]+)/force-close/?$', OpsTaskForceCloseView.as_view(), name='ops-task-force-close'),
180182

181183
# OpenAPI schema (support both with and without trailing slash)

app/api/views/ops.py

Lines changed: 58 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,34 @@
3434
logger = logging.getLogger('trawlr.api.ops')
3535

3636

37+
def _slim_queue_depth(q):
38+
"""
39+
Return just the fields the frontend renders for a queue depth row, or
40+
None if the queue doesn't exist on the broker. Avoids leaking the full
41+
RabbitMQ management payload over the API.
42+
"""
43+
if not q:
44+
return None
45+
return {
46+
'messages_ready': q.get('messages_ready', 0),
47+
'messages_unacknowledged': q.get('messages_unacknowledged', 0),
48+
}
49+
50+
51+
def _slim_queue_rows(rows):
52+
"""Reshape `get_queues_panel_context`'s queue rows into the API contract."""
53+
out = []
54+
for row in rows:
55+
out.append({
56+
'name': row['name'],
57+
'main': _slim_queue_depth(row.get('main')),
58+
'dq': _slim_queue_depth(row.get('dq')),
59+
'xq': _slim_queue_depth(row.get('xq')),
60+
'consumers': row.get('consumers', []),
61+
})
62+
return out
63+
64+
3765
class OpsQueuesView(APIView):
3866
"""
3967
GET /api/v1/ops/queues
@@ -47,7 +75,7 @@ class OpsQueuesView(APIView):
4775
def get(self, request):
4876
ctx = get_queues_panel_context()
4977
return Response({
50-
'queues': ctx['queue_rows'],
78+
'queues': _slim_queue_rows(ctx['queue_rows']),
5179
'taskSummary': ctx['task_summary'],
5280
'orphans': ctx['orphans'],
5381
'error': ctx.get('queues_error'),
@@ -130,6 +158,35 @@ def post(self, request, name):
130158
})
131159

132160

161+
class OpsConnectionCloseView(APIView):
162+
"""
163+
POST /api/v1/ops/connections/close
164+
Body: {"name": "<connection name>"}
165+
166+
Force-close a single RabbitMQ consumer connection (worker will
167+
auto-reconnect). Mirrors the v1 web `ops:connection_close` action.
168+
"""
169+
170+
permission_classes = [IsAdminUser]
171+
172+
def post(self, request):
173+
conn_name = (request.data.get('name') or '').strip()
174+
if not conn_name:
175+
return Response(
176+
{'error': 'name is required'},
177+
status=status.HTTP_400_BAD_REQUEST,
178+
)
179+
ok = rabbitmq.close_connection(conn_name, reason='closed via API')
180+
return Response(
181+
{'closed': bool(ok), 'name': conn_name},
182+
status=(
183+
status.HTTP_200_OK
184+
if ok
185+
else status.HTTP_502_BAD_GATEWAY
186+
),
187+
)
188+
189+
133190
class OpsTaskForceCloseView(APIView):
134191
"""
135192
POST /api/v1/ops/tasks/{id}/force-close

0 commit comments

Comments
 (0)