Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 11 additions & 5 deletions cpp/docs/grpc-server-architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -291,11 +291,17 @@ If a worker process crashes:

### Graceful Shutdown

On SIGINT/SIGTERM:
1. Set `shm_ctrl->shutdown_requested = true`
2. Workers finish current job and exit
3. Main process waits for workers
4. Cleanup shared memory segments
On SIGINT/SIGTERM (delivered to a dedicated `sigwait` thread — async `signal()`
handlers are unreliable once gRPC/CUDA threads mask signals):
1. Set `keep_running = false` and `shm_ctrl->shutdown_requested = true`
2. Mark all queued/running jobs `CANCELLED` and wake any `WaitForCompletion` waiters
3. `SIGKILL` all worker processes (mid-solve workers do not poll the shutdown flag)
4. Close server-side worker pipes so background threads blocked on pipe I/O unblock
5. Shut down the gRPC server with a short deadline so lingering RPCs cannot block exit
6. Join background threads, wait briefly for workers (`waitpid` with a grace period), and clean up shared memory
7. A 3s watchdog calls `_exit(0)` if the clean path wedges (e.g. GPU driver D-state)

Workers ignore SIGINT/SIGTERM so only the parent process owns shutdown; the parent always force-kills workers rather than waiting for the current solve to finish. Mid-CUDA workers can sit in uninterruptible D-state after `SIGKILL`, so `waitpid` is bounded — stragglers are abandoned for init/the test harness to reap.

### Job Cancellation

Expand Down
74 changes: 67 additions & 7 deletions cpp/src/grpc/server/grpc_server_main.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
#include <argparse/argparse.hpp>
#include <cuopt/version_config.hpp>

#include <pthread.h>

// Defined in grpc_service_impl.cpp
std::unique_ptr<grpc::Service> create_cuopt_grpc_service();

Expand Down Expand Up @@ -184,8 +186,19 @@ int main(int argc, char** argv)
SERVER_LOG_INFO("Port: %d", config.port);
SERVER_LOG_INFO("Workers: %d", config.num_workers);

signal(SIGINT, signal_handler);
signal(SIGTERM, signal_handler);
// Block shutdown signals in this thread before spawning workers / gRPC /
// background threads so they inherit the mask. A dedicated sigwait thread
// (started later) is then the only place SIGINT/SIGTERM are consumed.
// Async signal() handlers are unreliable once libraries mask signals.
sigset_t shutdown_sigset;
sigemptyset(&shutdown_sigset);
sigaddset(&shutdown_sigset, SIGINT);
sigaddset(&shutdown_sigset, SIGTERM);
int mask_rc = pthread_sigmask(SIG_BLOCK, &shutdown_sigset, nullptr);
if (mask_rc != 0) {
SERVER_LOG_ERROR("[Server] Failed to block shutdown signals: %s", strerror(mask_rc));
return 1;
}

ensure_log_dir_exists();

Expand Down Expand Up @@ -263,7 +276,17 @@ int main(int argc, char** argv)
log_worker_gpu_layout();
spawn_workers();

cuopt_expects(!worker_pids.empty(), error_type_t::RuntimeError, "No workers started");
{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
bool any_worker = false;
for (pid_t pid : worker_pids) {
if (pid > 0) {
any_worker = true;
break;
}
}
cuopt_expects(any_worker, error_type_t::RuntimeError, "No workers started");
}

} catch (const cuopt::logic_error& e) {
SERVER_LOG_ERROR("[Server] %s", format_cuopt_error(e));
Expand All @@ -288,6 +311,11 @@ int main(int argc, char** argv)
auto shutdown_all = [&]() {
keep_running = false;
shm_ctrl->shutdown_requested = true;
cancel_all_active_jobs_for_shutdown();
kill_all_workers();
// Close our pipe ends so any background thread blocked in write/read
// returns instead of delaying join past Ctrl-C.
close_all_server_worker_pipes();
result_cv.notify_all();

if (result_thread.joinable()) result_thread.join();
Expand Down Expand Up @@ -324,11 +352,43 @@ int main(int argc, char** argv)
server_max_message_bytes() / kMiB);
SERVER_LOG_INFO("[gRPC Server] Press Ctrl+C to shutdown");

std::thread shutdown_thread([&server]() {
while (keep_running.load()) {
std::this_thread::sleep_for(std::chrono::milliseconds(100));
// Dedicated signal thread: sigwait is reliable in multithreaded processes
// where gRPC/CUDA may mask signals on their threads. On SIGINT/SIGTERM we
// cancel jobs, kill workers, close pipes, and shut down gRPC. A watchdog
// forces process exit if cleanup hangs (e.g. GPU driver D-state).
sigset_t shutdown_sigset;
sigemptyset(&shutdown_sigset);
sigaddset(&shutdown_sigset, SIGINT);
sigaddset(&shutdown_sigset, SIGTERM);

std::thread shutdown_thread([&server, shutdown_sigset]() mutable {
int sig = 0;
int rc = sigwait(&shutdown_sigset, &sig);
if (rc != 0) {
SERVER_LOG_ERROR("[Server] sigwait failed: %s", strerror(rc));
sig = SIGTERM;
}

SERVER_LOG_INFO("[Server] Shutdown signal %d received; cancelling jobs and killing workers",
sig);
keep_running = false;
if (shm_ctrl) { shm_ctrl->shutdown_requested = true; }

// If joins / waitpid / gRPC Shutdown wedge, still die promptly so Ctrl-C
// and the integration test cannot hang indefinitely.
std::thread([]() {
std::this_thread::sleep_for(std::chrono::seconds(3));
SERVER_LOG_ERROR("[Server] Shutdown watchdog expired; forcing exit");
_exit(0);
}).detach();

cancel_all_active_jobs_for_shutdown();
kill_all_workers();
close_all_server_worker_pipes();
if (server) {
auto deadline = std::chrono::system_clock::now() + std::chrono::seconds(1);
server->Shutdown(deadline);
}
if (server) { server->Shutdown(); }
});

server->Wait();
Expand Down
68 changes: 44 additions & 24 deletions cpp/src/grpc/server/grpc_server_threads.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -14,44 +14,64 @@ void worker_monitor_thread()
SERVER_LOG_INFO("[Server] Worker monitor thread started");

while (keep_running) {
for (size_t i = 0; i < worker_pids.size(); ++i) {
pid_t pid = worker_pids[i];
if (pid <= 0) continue;
// Snapshot which slots need attention under the pid-list lock, then do
// mark/respawn work without holding it across fork().
struct DeadWorker {
size_t index;
pid_t pid;
bool was_clean_shutdown_exit;
};
std::vector<DeadWorker> dead;

int status;
pid_t result = waitpid(pid, &status, WNOHANG);
{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
for (size_t i = 0; i < worker_pids.size(); ++i) {
pid_t pid = worker_pids[i];
if (pid <= 0) continue;

int status = 0;
pid_t wait_ret = waitpid(pid, &status, WNOHANG);
if (wait_ret != pid) continue;

if (result == pid) {
int exit_code = WIFEXITED(status) ? WEXITSTATUS(status) : -1;
bool signaled = WIFSIGNALED(status);
int signal_num = signaled ? WTERMSIG(status) : 0;
int exit_code = WIFEXITED(status) ? WEXITSTATUS(status) : -1;
bool signaled = WIFSIGNALED(status);
int signal_num = signaled ? WTERMSIG(status) : 0;
bool was_clean_shutdown_exit = false;

if (signaled) {
SERVER_LOG_ERROR("[Server] Worker %d killed by signal %d", pid, signal_num);
} else if (exit_code != 0) {
SERVER_LOG_ERROR("[Server] Worker %d exited with code %d", pid, exit_code);
} else if (shm_ctrl && shm_ctrl->shutdown_requested) {
was_clean_shutdown_exit = true;
} else {
if (shm_ctrl && shm_ctrl->shutdown_requested) {
worker_pids[i] = 0;
continue;
}
SERVER_LOG_ERROR("[Server] Worker %d exited unexpectedly", pid);
}

mark_worker_jobs_failed(pid);
worker_pids[i] = 0;
dead.push_back({i, pid, was_clean_shutdown_exit});
}
}

if (keep_running && shm_ctrl && !shm_ctrl->shutdown_requested) {
pid_t new_pid = spawn_single_worker(static_cast<int>(i));
if (new_pid > 0) {
worker_pids[i] = new_pid;
SERVER_LOG_INFO("[Server] Restarted worker %zu with PID %d", i, new_pid);
} else {
worker_pids[i] = 0;
}
} else {
worker_pids[i] = 0;
for (const auto& dw : dead) {
if (dw.was_clean_shutdown_exit) continue;

mark_worker_jobs_failed(dw.pid);

if (!(keep_running && shm_ctrl && !shm_ctrl->shutdown_requested)) { continue; }

pid_t new_pid = spawn_single_worker(static_cast<int>(dw.index));
{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
if (dw.index < worker_pids.size() && worker_pids[dw.index] == 0) {
worker_pids[dw.index] = (new_pid > 0) ? new_pid : 0;
}
}
if (new_pid > 0) {
SERVER_LOG_INFO("[Server] Restarted worker %zu with PID %d", dw.index, new_pid);
} else {
SERVER_LOG_ERROR("[Server] Failed to restart worker %zu", dw.index);
}
Comment thread
tmckayus marked this conversation as resolved.
}

std::this_thread::sleep_for(std::chrono::milliseconds(100));
Expand Down
13 changes: 6 additions & 7 deletions cpp/src/grpc/server/grpc_server_types.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,7 @@ inline ResultQueueEntry* result_queue = nullptr;
inline SharedMemoryControl* shm_ctrl = nullptr;

inline std::vector<pid_t> worker_pids;
inline std::mutex worker_pids_mutex;

inline ServerConfig config;

Expand Down Expand Up @@ -330,13 +331,8 @@ inline std::string read_file_to_string(const std::string& path)
// Signal handling
// =============================================================================

inline void signal_handler(int signal)
{
if (signal == SIGINT || signal == SIGTERM) {
keep_running = false;
if (shm_ctrl) { shm_ctrl->shutdown_requested = true; }
}
}
// SIGINT/SIGTERM are handled via sigwait on a dedicated thread (see main).
// Using signal() handlers is unreliable once gRPC/CUDA threads mask signals.

// =============================================================================
// Forward declarations
Expand All @@ -350,7 +346,10 @@ void cleanup_shared_memory();
void log_worker_gpu_layout();
bool init_worker_cuda_environment(int worker_id);
void spawn_workers();
void kill_all_workers();
void close_all_server_worker_pipes();
void wait_for_workers();
void cancel_all_active_jobs_for_shutdown();
void worker_monitor_thread();
void result_retrieval_thread();
void incumbent_retrieval_thread();
Expand Down
6 changes: 6 additions & 0 deletions cpp/src/grpc/server/grpc_worker.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -602,6 +602,12 @@ void worker_process(int worker_id)
{
SERVER_LOG_INFO("[Worker %d] Started (PID: %d)", worker_id, getpid());

// Parent owns SIGINT/SIGTERM shutdown. Ignoring here prevents the inherited
// soft handler from leaving mid-solve workers alive after Ctrl-C while the
// parent waits on them.
signal(SIGINT, SIG_IGN);
signal(SIGTERM, SIG_IGN);

if (!init_worker_cuda_environment(worker_id)) {
SERVER_LOG_ERROR("[Worker %d] CUDA environment initialization failed; exiting", worker_id);
_exit(1);
Expand Down
112 changes: 105 additions & 7 deletions cpp/src/grpc/server/grpc_worker_infra.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -201,21 +201,119 @@ pid_t spawn_worker(int worker_id, bool is_replacement)

void spawn_workers()
{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
// Index i is worker_id: keep failed startups as 0 so monitor/respawn and
// pipe tables stay aligned even when some initial forks fail.
worker_pids.assign(static_cast<size_t>(config.num_workers), 0);
for (int i = 0; i < config.num_workers; ++i) {
pid_t pid = spawn_worker(i, false);
if (pid < 0) { continue; }
worker_pids.push_back(pid);
if (pid > 0) { worker_pids[static_cast<size_t>(i)] = pid; }
}
}

void wait_for_workers()
void kill_all_workers()
{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
for (pid_t pid : worker_pids) {
if (pid <= 0) continue;
int status;
while (waitpid(pid, &status, 0) < 0 && errno == EINTR) {}
if (pid > 0) { kill(pid, SIGKILL); }
}
}
Comment thread
tmckayus marked this conversation as resolved.

void close_all_server_worker_pipes()
{
std::lock_guard<std::mutex> lock(worker_pipes_mutex);
for (auto& wp : worker_pipes) {
close_all_worker_pipes(wp);
}
}

void cancel_all_active_jobs_for_shutdown()
{
if (job_queue != nullptr) {
for (size_t i = 0; i < MAX_JOBS; ++i) {
if (job_queue[i].ready.load(std::memory_order_acquire)) {
job_queue[i].cancelled.store(true, std::memory_order_release);
}
}
}

{
std::lock_guard<std::mutex> lock(tracker_mutex);
for (auto& [job_id, info] : job_tracker) {
(void)job_id;
if (info.status == JobStatus::QUEUED || info.status == JobStatus::PROCESSING) {
info.status = JobStatus::CANCELLED;
info.error_message = "Server shutting down";
}
}
}

{
std::lock_guard<std::mutex> wlock(waiters_mutex);
for (auto& [job_id, waiter] : waiting_threads) {
(void)job_id;
{
std::lock_guard<std::mutex> waiter_lock(waiter->mutex);
waiter->error_message = "Server shutting down";
waiter->success = false;
waiter->ready = true;
}
waiter->cv.notify_all();
}
waiting_threads.clear();
}

result_cv.notify_all();
}

void wait_for_workers()
{
// Mid-solve workers only check shutdown_requested between jobs, so force-kill
// before waitpid or Ctrl-C / SIGTERM can hang until the current solve ends.
// A worker killed mid-CUDA can sit in uninterruptible D-state while the GPU
// driver tears down; a blocking waitpid would then hang the whole server
// (and ignore SIGTERM because we retry on EINTR). Bound the wait and abandon
// stragglers — the test harness / init reaps them.
kill_all_workers();

constexpr auto kShutdownWait = std::chrono::seconds(2);
auto deadline = std::chrono::steady_clock::now() + kShutdownWait;
while (std::chrono::steady_clock::now() < deadline) {
bool any_alive = false;
{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
for (pid_t& pid : worker_pids) {
if (pid <= 0) continue;
int status = 0;
pid_t reaped = waitpid(pid, &status, WNOHANG);
if (reaped == pid || (reaped < 0 && errno == ECHILD)) {
pid = 0;
} else if (reaped < 0 && errno == EINTR) {
any_alive = true;
} else {
any_alive = true;
}
}
if (!any_alive) {
worker_pids.clear();
return;
}
}
std::this_thread::sleep_for(std::chrono::milliseconds(50));
}

{
std::lock_guard<std::mutex> lock(worker_pids_mutex);
for (pid_t pid : worker_pids) {
if (pid > 0) {
kill(pid, SIGKILL);
SERVER_LOG_WARN(
"[Server] Worker pid %d did not exit within shutdown grace period; abandoning",
static_cast<int>(pid));
}
}
worker_pids.clear();
}
worker_pids.clear();
}

pid_t spawn_single_worker(int worker_id) { return spawn_worker(worker_id, true); }
Expand Down
Loading
Loading