| 12 | |
| 13 | |
| 14 | async def cancel_queued_request( |
| 15 | req_id: str, request_queue: Queue, logger: logging.Logger |
| 16 | ) -> bool: |
| 17 | set_request_id(req_id) |
| 18 | items_to_requeue = [] |
| 19 | found = False |
| 20 | try: |
| 21 | while not request_queue.empty(): |
| 22 | item = request_queue.get_nowait() |
| 23 | if item.get("req_id") == req_id: |
| 24 | logger.info("Found request in queue, marking as cancelled.") |
| 25 | item["cancelled"] = True |
| 26 | if (future := item.get("result_future")) and not future.done(): |
| 27 | future.set_exception(client_cancelled(req_id)) |
| 28 | found = True |
| 29 | items_to_requeue.append(item) |
| 30 | finally: |
| 31 | for item in items_to_requeue: |
| 32 | await request_queue.put(item) |
| 33 | return found |
| 34 | |
| 35 | |
| 36 | async def cancel_request( |