| 228 | } |
| 229 | |
| 230 | void EventLoop::loop() |
| 231 | { |
| 232 | assert(!g_thread_context.loop_thread); |
| 233 | g_thread_context.loop_thread = true; |
| 234 | KJ_DEFER(g_thread_context.loop_thread = false); |
| 235 | |
| 236 | { |
| 237 | const Lock lock(m_mutex); |
| 238 | assert(!m_async_fns); |
| 239 | m_async_fns.emplace(); |
| 240 | } |
| 241 | |
| 242 | kj::Own<kj::AsyncIoStream> wait_stream{ |
| 243 | m_io_context.lowLevelProvider->wrapSocketFd(m_wait_fd, kj::LowLevelAsyncIoProvider::TAKE_OWNERSHIP)}; |
| 244 | int post_fd{m_post_fd}; |
| 245 | char buffer = 0; |
| 246 | for (;;) { |
| 247 | const size_t read_bytes = wait_stream->read(&buffer, 0, 1).wait(m_io_context.waitScope); |
| 248 | if (read_bytes != 1) throw std::logic_error("EventLoop wait_stream closed unexpectedly"); |
| 249 | Lock lock(m_mutex); |
| 250 | if (m_post_fn) { |
| 251 | // m_post_fn throwing is never expected. If it does happen, the caller |
| 252 | // of EventLoop::post() will return without any indication of failure, |
| 253 | // which will likely cause other bugs. Log the error and continue. |
| 254 | KJ_IF_MAYBE(exception, kj::runCatchingExceptions([&]() MP_REQUIRES(m_mutex) { Unlock(lock, *m_post_fn); })) { |
| 255 | MP_LOG(*this, Log::Error) << "EventLoop: m_post_fn threw: " << kj::str(*exception).cStr(); |
| 256 | } |
| 257 | m_post_fn = nullptr; |
| 258 | m_cv.notify_all(); |
| 259 | } else if (done()) { |
| 260 | // Intentionally do not break if m_post_fn was set, even if done() |
| 261 | // would return true, to ensure that the EventLoopRef write(post_fd) |
| 262 | // call always succeeds and the loop does not exit between the time |
| 263 | // that the done condition is set and the write call is made. |
| 264 | break; |
| 265 | } |
| 266 | } |
| 267 | MP_LOG(*this, Log::Info) << "EventLoop::loop done, cancelling event listeners."; |
| 268 | m_task_set.reset(); |
| 269 | MP_LOG(*this, Log::Info) << "EventLoop::loop bye."; |
| 270 | wait_stream = nullptr; |
| 271 | KJ_SYSCALL(::close(post_fd)); |
| 272 | const Lock lock(m_mutex); |
| 273 | m_wait_fd = -1; |
| 274 | m_post_fd = -1; |
| 275 | m_async_fns.reset(); |
| 276 | m_cv.notify_all(); |
| 277 | } |
| 278 | |
| 279 | void EventLoop::post(kj::Function<void()> fn) |
| 280 | { |