| 220 | } |
| 221 | |
| 222 | bool Wait(TEvent& ev, const TInstant deadline_ = TInstant::Max()) override { |
| 223 | while (!Interrupt_) { |
| 224 | TInstant deadline = deadline_; |
| 225 | const TInstant now = TInstant::Now(); |
| 226 | if (deadline != TInstant::Max() && now >= deadline) { |
| 227 | break; |
| 228 | } |
| 229 | |
| 230 | { //process jobs queue (requests/responses info) |
| 231 | TAutoPtr<IJob> j; |
| 232 | while (JQ_.Dequeue(&j)) { |
| 233 | if (j->Process(ev)) { |
| 234 | return true; |
| 235 | } |
| 236 | } |
| 237 | } |
| 238 | |
| 239 | if (!RS_.Empty()) { |
| 240 | TRequestSupervisor* nearRS = &*RS_.Begin(); |
| 241 | if (nearRS->Deadline() <= now) { |
| 242 | if (!nearRS->MarkAsHandled()) { |
| 243 | //race with notify, - now in queue must exist response job for this request |
| 244 | continue; |
| 245 | } |
| 246 | ev.Type = TEvent::Timeout; |
| 247 | nearRS->FillEvent(ev); |
| 248 | nearRS->ResetRequest(); |
| 249 | nearRS->UnLink(); |
| 250 | return true; |
| 251 | } |
| 252 | deadline = Min(nearRS->Deadline(), deadline); |
| 253 | } |
| 254 | |
| 255 | if (SetNearDeadline(deadline)) { |
| 256 | continue; //update deadline to more far time, so need re-check queue for avoiding race |
| 257 | } |
| 258 | |
| 259 | E_.WaitD(deadline); |
| 260 | } |
| 261 | Interrupt_ = false; |
| 262 | return false; |
| 263 | } |
| 264 | |
| 265 | void Interrupt() override { |
| 266 | Interrupt_ = true; |
no test coverage detected