Process the next pending task using the 8:4:2 priority drain ratio. Advances `self.drain_cycle` by one slot and returns `true` if a task was processed.
(&mut self)
| 34 | /// Advances `self.drain_cycle` by one slot and returns `true` if a task |
| 35 | /// was processed. |
| 36 | pub fn poll_one(&mut self) -> bool { |
| 37 | let Some(qt) = self.task_queue.pop_next(&mut self.drain_cycle) else { |
| 38 | return false; |
| 39 | }; |
| 40 | |
| 41 | // Record IO wait from enqueue to execution start. |
| 42 | let wait_ns = qt.enqueued_at.elapsed().as_nanos() as u64; |
| 43 | use crate::bridge::envelope::Priority; |
| 44 | let tier = match qt.task.request.priority { |
| 45 | Priority::Background | Priority::Normal => TIER_LOW, |
| 46 | Priority::High => TIER_HIGH, |
| 47 | Priority::Critical => TIER_CRITICAL, |
| 48 | }; |
| 49 | self.io_metrics.record_wait(tier, wait_ns); |
| 50 | |
| 51 | let mut task = qt.task; |
| 52 | |
| 53 | if let Some(key) = task.request.idempotency_key |
| 54 | && let Some(&succeeded) = self.idempotency_cache.get(&key) |
| 55 | { |
| 56 | let response = if succeeded { |
| 57 | self.response_ok(&task) |
| 58 | } else { |
| 59 | self.response_error(&task, ErrorCode::DuplicateWrite) |
| 60 | }; |
| 61 | if let Err(e) = self |
| 62 | .response_tx |
| 63 | .try_push(BridgeResponse { inner: response }) |
| 64 | { |
| 65 | warn!(core = self.core_id, error = %e, "failed to send idempotent response"); |
| 66 | } |
| 67 | return true; |
| 68 | } |
| 69 | |
| 70 | let response = if task.is_expired() { |
| 71 | task.state = TaskState::Failed; |
| 72 | Response { |
| 73 | request_id: task.request_id(), |
| 74 | status: Status::Error, |
| 75 | attempt: 1, |
| 76 | partial: false, |
| 77 | payload: Payload::empty(), |
| 78 | watermark_lsn: self.watermark, |
| 79 | error_code: Some(ErrorCode::DeadlineExceeded), |
| 80 | } |
| 81 | } else { |
| 82 | task.state = TaskState::Running; |
| 83 | let resp = self.execute(&task); |
| 84 | task.state = TaskState::Completed; |
| 85 | resp |
| 86 | }; |
| 87 | |
| 88 | if let Some(key) = task.request.idempotency_key { |
| 89 | let succeeded = response.status == Status::Ok; |
| 90 | if self.idempotency_cache.len() >= 16_384 |
| 91 | && let Some(oldest_key) = self.idempotency_order.pop_front() |
| 92 | { |
| 93 | self.idempotency_cache.remove(&oldest_key); |
no test coverage detected