| 58 | |
| 59 | |
| 60 | async def get_queue_status( |
| 61 | request_queue: Queue = Depends(get_request_queue), |
| 62 | processing_lock: Lock = Depends(get_processing_lock), |
| 63 | ): |
| 64 | # Extract all items temporarily to inspect queue contents |
| 65 | queue_items = [] |
| 66 | try: |
| 67 | while not request_queue.empty(): |
| 68 | item = request_queue.get_nowait() |
| 69 | queue_items.append(item) |
| 70 | except Exception: |
| 71 | pass |
| 72 | finally: |
| 73 | # Put all items back in original order |
| 74 | for item in queue_items: |
| 75 | await request_queue.put(item) |
| 76 | |
| 77 | queue_length = len(queue_items) |
| 78 | |
| 79 | return JSONResponse( |
| 80 | content={ |
| 81 | "queue_length": queue_length, |
| 82 | "is_processing_locked": processing_lock.locked(), |
| 83 | "items": sorted( |
| 84 | [ |
| 85 | { |
| 86 | "req_id": item.get("req_id", "unknown"), |
| 87 | "enqueue_time": item.get("enqueue_time", 0), |
| 88 | "wait_time_seconds": round( |
| 89 | time.time() - item.get("enqueue_time", 0), 2 |
| 90 | ), |
| 91 | "is_streaming": item.get("request_data").stream, |
| 92 | "cancelled": item.get("cancelled", False), |
| 93 | } |
| 94 | for item in queue_items |
| 95 | ], |
| 96 | key=lambda x: x.get("enqueue_time", 0), |
| 97 | ), |
| 98 | } |
| 99 | ) |