| 82 | } |
| 83 | |
| 84 | def _update_queue_metrics(self, queue_sizes: dict[str, int], transfer_sizes: Optional[dict[str, int]] = None): |
| 85 | merged_sizes = {k: int(v) for k, v in queue_sizes.items()} |
| 86 | if transfer_sizes is not None: |
| 87 | for key, value in transfer_sizes.items(): |
| 88 | merged_sizes[f"transfer_{key}"] = int(value) |
| 89 | total_pending = sum(max(v, 0) for v in merged_sizes.values()) |
| 90 | with self._queue_metrics_lock: |
| 91 | self._queue_metrics = { |
| 92 | "queue_sizes": merged_sizes, |
| 93 | "queue_total_pending": total_pending, |
| 94 | "all_queues_empty": total_pending == 0, |
| 95 | } |
| 96 | |
| 97 | def _ensure_phase2_request_buffer(self) -> bool: |
| 98 | if self._phase2_rdma_buffer is not None: |