| 418 | ProxyServer<ThreadMap>::ProxyServer(Connection& connection) : m_connection(connection) {} |
| 419 | |
| 420 | kj::Promise<void> ProxyServer<ThreadMap>::makePool(MakePoolContext context) |
| 421 | { |
| 422 | if (!m_connection.m_thread_pool.empty()) { |
| 423 | throw std::runtime_error("makePool called on connection with existing pool"); |
| 424 | } |
| 425 | EventLoop& loop{*m_connection.m_loop}; |
| 426 | const uint32_t count = context.getParams().getCount(); |
| 427 | for (uint32_t i = 0; i < count; ++i) { |
| 428 | const std::string thread_name = "pool/" + std::to_string(i); |
| 429 | std::promise<ThreadContext*> thread_context; |
| 430 | std::thread thread([&loop, &thread_context, thread_name]() { |
| 431 | g_thread_context.thread_name = ThreadName(loop.m_exe_name) + " (" + thread_name + ")"; |
| 432 | g_thread_context.waiter = std::make_unique<Waiter>(); |
| 433 | Lock lock(g_thread_context.waiter->m_mutex); |
| 434 | thread_context.set_value(&g_thread_context); |
| 435 | g_thread_context.waiter->wait(lock, [] { return !g_thread_context.waiter; }); |
| 436 | }); |
| 437 | auto thread_server = kj::heap<ProxyServer<Thread>>(m_connection, *thread_context.get_future().get(), std::move(thread)); |
| 438 | m_connection.m_thread_pool.push_back({m_connection.m_threads.add(kj::mv(thread_server))}); |
| 439 | } |
| 440 | return kj::READY_NOW; |
| 441 | } |
| 442 | |
| 443 | kj::Promise<void> ProxyServer<ThreadMap>::makeThread(MakeThreadContext context) |
| 444 | { |