From dfeda4dc6f50b056b56c1928a6d68cf4ca8bc4b1 Mon Sep 17 00:00:00 2001 From: ykiko Date: Sun, 31 May 2026 23:48:40 +0800 Subject: [PATCH 01/10] fix(server): make shutdown sanitizer-clean --- src/server/service/agent_client.cpp | 3 +- src/server/service/lsp_client.cpp | 1 + src/server/service/master_server.cpp | 76 ++++++++++------ src/server/worker/worker_pool.cpp | 7 +- tests/conftest.py | 102 ++++++++++++++++------ tests/integration/agentic/test_agentic.py | 61 +++++++------ 6 files changed, 166 insertions(+), 84 deletions(-) diff --git a/src/server/service/agent_client.cpp b/src/server/service/agent_client.cpp index 69c7a81b2..d51b610e1 100644 --- a/src/server/service/agent_client.cpp +++ b/src/server/service/agent_client.cpp @@ -778,9 +778,10 @@ AgentClient::AgentClient(MasterServer& server, kota::ipc::JsonPeer& peer) : co_return result; }); - peer.on_notification([&srv](const ShutdownParams&) { + peer.on_notification([this, &srv](const ShutdownParams&) { LOG_INFO("agentic/shutdown received, shutting down"); srv.schedule_shutdown(); + this->peer.close(); }); } diff --git a/src/server/service/lsp_client.cpp b/src/server/service/lsp_client.cpp index 1607b60a7..b7bbc369c 100644 --- a/src/server/service/lsp_client.cpp +++ b/src/server/service/lsp_client.cpp @@ -156,6 +156,7 @@ LSPClient::LSPClient(MasterServer& server, kota::ipc::JsonPeer& peer) : server(s peer.on_notification([this]([[maybe_unused]] const protocol::ExitParams& params) { LOG_INFO("Exit notification received"); this->server.schedule_shutdown(); + this->peer.close(); }); peer.on_notification([this](const protocol::DidOpenTextDocumentParams& params) { diff --git a/src/server/service/master_server.cpp b/src/server/service/master_server.cpp index 9e8b9010c..656a3245d 100644 --- a/src/server/service/master_server.cpp +++ b/src/server/service/master_server.cpp @@ -212,10 +212,10 @@ void MasterServer::schedule_shutdown() { workspace.save_cache(); shutdown_event.set(); - loop.schedule([this]() -> kota::task<> { - co_await kota::when_all(indexer.stop(), compiler.stop(), pool.stop()); - loop.stop(); - }()); + loop.schedule([](MasterServer& server) -> kota::task<> { + co_await kota::when_all(server.indexer.stop(), server.compiler.stop(), server.pool.stop()); + server.loop.stop(); + }(*this)); } void MasterServer::load_workspace() { @@ -355,35 +355,53 @@ static kota::task<> accept_connections(MasterServer& server, std::list& connections) { auto& loop = kota::event_loop::current(); kota::task_group<> connection_group(loop); - bool lsp_registered = false; - while(true) { - auto conn = co_await acceptor.accept(); - if(!conn.has_value()) - break; + co_await kota::when_all( + [](MasterServer& server, + kota::tcp::acceptor& acceptor, + bool register_lsp, + std::list& connections, + kota::task_group<>& connection_group) -> kota::task<> { + auto& loop = kota::event_loop::current(); + bool lsp_registered = false; - LOG_INFO("Client connected"); + while(true) { + auto conn = co_await acceptor.accept(); + if(!conn.has_value()) + break; - auto transport = std::make_unique(std::move(*conn)); - auto peer = std::make_unique(loop, std::move(transport)); + LOG_INFO("Client connected"); - std::unique_ptr lsp; - if(register_lsp && !lsp_registered) { - lsp = std::make_unique(server, *peer); - lsp_registered = true; - } - auto agent = std::make_unique(server, *peer); + auto transport = std::make_unique(std::move(*conn)); + auto peer = std::make_unique(loop, std::move(transport)); - auto* peer_ptr = peer.get(); - auto it = connections.emplace(connections.end(), - Connection{ - .peer = std::move(peer), - .lsp_client = std::move(lsp), - .agent_client = std::move(agent), - }); + std::unique_ptr lsp; + if(register_lsp && !lsp_registered) { + lsp = std::make_unique(server, *peer); + lsp_registered = true; + } + auto agent = std::make_unique(server, *peer); - connection_group.spawn(run_connection(peer_ptr, connections, it)); - } + auto* peer_ptr = peer.get(); + auto it = connections.emplace(connections.end(), + Connection{ + .peer = std::move(peer), + .lsp_client = std::move(lsp), + .agent_client = std::move(agent), + }); + + connection_group.spawn(run_connection(peer_ptr, connections, it)); + } + }(server, acceptor, register_lsp, connections, connection_group), + [](MasterServer& server, + kota::tcp::acceptor& acceptor, + std::list& connections) -> kota::task<> { + co_await server.get_shutdown_event().wait(); + acceptor.stop(); + for(auto& conn: connections) { + conn.peer->close(); + } + }(server, acceptor, connections)); co_await connection_group.join(); } @@ -411,6 +429,10 @@ int run_server_mode(const ServerOptions& opts) { kota::ipc::JsonPeer lsp_peer(loop, std::move(final_transport)); LSPClient lsp_client(server, lsp_peer); + loop.schedule([](MasterServer& server, kota::ipc::JsonPeer& peer) -> kota::task<> { + co_await server.get_shutdown_event().wait(); + peer.close(); + }(server, lsp_peer)); if(opts.port > 0) { auto acceptor = kota::tcp::listen(opts.host, opts.port, {}, loop); diff --git a/src/server/worker/worker_pool.cpp b/src/server/worker/worker_pool.cpp index 2253e1b82..07adc35d2 100644 --- a/src/server/worker/worker_pool.cpp +++ b/src/server/worker/worker_pool.cpp @@ -116,6 +116,9 @@ bool WorkerPool::start(const WorkerPoolOptions& options) { options_ = options; log_dir_ = options.log_dir; + stateless_workers.reserve(options.stateless_count); + stateful_workers.reserve(options.stateful_count); + for(std::uint32_t i = 0; i < options.stateless_count; ++i) { if(!spawn_worker(options.self_path, false, 0)) { return false; @@ -229,10 +232,10 @@ void WorkerPool::clear_owner(std::size_t worker_index) { kota::task<> WorkerPool::monitor_worker(std::size_t index, bool stateful) { auto& workers = stateful ? stateful_workers : stateless_workers; - auto& w = workers[index]; auto name = std::string(stateful ? "SF-" : "SL-") + std::to_string(index); - auto result = co_await w.proc.wait(); + auto result = co_await workers[index].proc.wait(); + auto& w = workers[index]; w.alive = false; if(shutting_down_) diff --git a/tests/conftest.py b/tests/conftest.py index 05890cc11..ecca7a0ad 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -185,44 +185,94 @@ async def make_client(executable: Path, workspace: Path) -> CliceClient: return c +SANITIZER_MARKERS = ( + "AddressSanitizer", + "LeakSanitizer", + "MemorySanitizer", + "ThreadSanitizer", + "UndefinedBehaviorSanitizer", + "==ERROR:", + "runtime error:", +) + + +def _server_stderr_excerpt(stderr_text: str) -> str: + interesting = [ + line + for line in stderr_text.splitlines() + if "[warn]" in line + or "[error]" in line + or "Sanitizer" in line + or "==ERROR:" in line + or "runtime error:" in line + ] + return "\n".join(interesting[-80:]) + + +async def assert_server_exited_cleanly(server, timeout: float = 3.0) -> None: + failures: list[str] = [] + + if server is None: + return + + if server.returncode is None: + try: + await asyncio.wait_for(server.wait(), timeout=timeout) + except asyncio.TimeoutError: + server.kill() + await server.wait() + failures.append(f"server did not exit within {timeout:g}s after shutdown") + + print(f"[server] exit code: {server.returncode}", flush=True) + + stderr_text = "" + if server.stderr: + try: + stderr_data = await asyncio.wait_for(server.stderr.read(), timeout=2.0) + stderr_text = stderr_data.decode("utf-8", errors="replace") + except Exception as exc: + failures.append(f"failed to collect server stderr: {exc!r}") + + for line in _server_stderr_excerpt(stderr_text).splitlines(): + print(f"[server] {line}", flush=True) + + if server.returncode != 0: + failures.append(f"server exited with code {server.returncode}") + + if any(marker in stderr_text for marker in SANITIZER_MARKERS): + failures.append("server stderr contains sanitizer/runtime error output") + + if failures: + excerpt = _server_stderr_excerpt(stderr_text) + if excerpt: + failures.append("server stderr excerpt:\n" + excerpt) + pytest.fail("\n".join(failures)) + + async def _shutdown_client(c: CliceClient) -> None: """Gracefully shut down a client, force-kill if needed.""" + server = getattr(c, "_server", None) + try: await asyncio.wait_for(c.shutdown_async(None), timeout=3.0) except Exception: pass - try: - c.exit(None) - except Exception: - pass - - await asyncio.sleep(0.3) - if hasattr(c, "_server") and c._server is not None and c._server.returncode is None: - c._server.kill() try: - server = getattr(c, "_server", None) - if server: - if server.returncode is not None: - print(f"[server] exit code: {server.returncode}", flush=True) - if server.stderr: - stderr_data = await asyncio.wait_for(server.stderr.read(), timeout=2.0) - if stderr_data: - for line in stderr_data.decode( - "utf-8", errors="replace" - ).splitlines(): - if "[warn]" in line or "[error]" in line or "Sanitizer" in line: - print(f"[server] {line}", flush=True) + c.exit(None) except Exception: pass try: - c._stop_event.set() - for task in c._async_tasks: - task.cancel() - await asyncio.sleep(0.1) - except Exception: - pass + await assert_server_exited_cleanly(server) + finally: + try: + c._stop_event.set() + for task in c._async_tasks: + task.cancel() + await asyncio.sleep(0.1) + except Exception: + pass shutdown_client = _shutdown_client # Public alias for multi-session tests diff --git a/tests/integration/agentic/test_agentic.py b/tests/integration/agentic/test_agentic.py index 549647a4b..75ad08157 100644 --- a/tests/integration/agentic/test_agentic.py +++ b/tests/integration/agentic/test_agentic.py @@ -530,7 +530,7 @@ async def test_rpc_impact_analysis_unknown(indexed_agentic, workspace): async def test_shutdown_during_indexing(executable, tmp_path): """Shutdown during active background indexing must exit cleanly.""" from tests.integration.utils.client import CliceClient - from tests.conftest import _find_free_port + from tests.conftest import _find_free_port, assert_server_exited_cleanly workspace = tmp_path / "ws" workspace.mkdir() @@ -560,33 +560,38 @@ async def test_shutdown_during_indexing(executable, tmp_path): c = CliceClient() await c.start_io(*cmd) - init_options = { - "project": { - "cache_dir": str(workspace / ".clice"), - "idle_timeout_ms": 0, - } - } - await c.initialize(workspace, initialization_options=init_options) - - # Give indexing a moment to start, then send shutdown - await asyncio.sleep(0.5) - - rpc = AgenticRpcClient(host, port) - body = json.dumps({"jsonrpc": "2.0", "method": "agentic/shutdown", "params": {}}) - rpc.sock.sendall(f"Content-Length: {len(body)}\r\n\r\n{body}".encode()) - rpc.sock.settimeout(5) try: - rpc.sock.recv(4096) - except (socket.timeout, OSError): - pass - rpc.sock.close() - - for _ in range(30): - if c._server.returncode is not None: - break + init_options = { + "project": { + "cache_dir": str(workspace / ".clice"), + "idle_timeout_ms": 0, + } + } + try: + await c.initialize(workspace, initialization_options=init_options) + except Exception: + if c._server.returncode is not None: + await assert_server_exited_cleanly(c._server, timeout=15.0) + raise + + # Give indexing a moment to start, then send shutdown await asyncio.sleep(0.5) - assert c._server.returncode is not None, "Server did not exit after shutdown" - assert c._server.returncode >= 0, ( - f"Server crashed with signal {-c._server.returncode}" - ) + rpc = AgenticRpcClient(host, port) + body = json.dumps( + {"jsonrpc": "2.0", "method": "agentic/shutdown", "params": {}} + ) + rpc.sock.sendall(f"Content-Length: {len(body)}\r\n\r\n{body}".encode()) + rpc.sock.settimeout(5) + try: + rpc.sock.recv(4096) + except (socket.timeout, OSError): + pass + rpc.sock.close() + + await assert_server_exited_cleanly(c._server, timeout=15.0) + finally: + c._stop_event.set() + for task in c._async_tasks: + task.cancel() + await asyncio.sleep(0.1) From 42a3d2097127b4bee6a6c76e26ea182c4349b745 Mon Sep 17 00:00:00 2001 From: ykiko Date: Sun, 31 May 2026 23:48:40 +0800 Subject: [PATCH 02/10] fix(server): make shutdown sanitizer-clean --- src/server/service/master_server.cpp | 63 +++++++++++++++++++++------- 1 file changed, 49 insertions(+), 14 deletions(-) diff --git a/src/server/service/master_server.cpp b/src/server/service/master_server.cpp index 656a3245d..7ed6cd653 100644 --- a/src/server/service/master_server.cpp +++ b/src/server/service/master_server.cpp @@ -110,14 +110,14 @@ void MasterServer::start_file_watcher() { if(workspace_root.empty()) return; - loop.schedule([this]() -> kota::task<> { - auto watcher = kota::fs_event::create(workspace_root, {}, loop); + loop.schedule([](MasterServer& server) -> kota::task<> { + auto watcher = kota::fs_event::create(server.workspace_root, {}, server.loop); if(!watcher) { - LOG_WARN("Failed to start file watcher for {}", workspace_root); + LOG_WARN("Failed to start file watcher for {}", server.workspace_root); co_return; } - LOG_INFO("File watcher started for {}", workspace_root); + LOG_INFO("File watcher started for {}", server.workspace_root); while(true) { auto changes = co_await watcher->next(); @@ -132,19 +132,19 @@ void MasterServer::start_file_watcher() { llvm::StringRef file(change.path); if(file.ends_with("compile_commands.json")) { LOG_INFO("CDB changed, reloading workspace"); - load_workspace(); + server.load_workspace(); continue; } if(file.ends_with(".cpp") || file.ends_with(".cc") || file.ends_with(".cxx") || file.ends_with(".c") || file.ends_with(".h") || file.ends_with(".hpp") || file.ends_with(".hxx") || file.ends_with(".cppm") || file.ends_with(".ixx")) { - auto path_id = workspace.path_pool.intern(file); - on_file_saved(path_id); + auto path_id = server.workspace.path_pool.intern(file); + server.on_file_saved(path_id); } } } - }()); + }(*this)); } Session* MasterServer::find_session(std::uint32_t path_id) { @@ -434,17 +434,45 @@ int run_server_mode(const ServerOptions& opts) { peer.close(); }(server, lsp_peer)); + kota::tcp::acceptor agent_acceptor; + bool has_agent_acceptor = false; + if(opts.port > 0) { auto acceptor = kota::tcp::listen(opts.host, opts.port, {}, loop); if(acceptor) { LOG_INFO("Agentic protocol listening on {}:{}", opts.host, opts.port); - loop.schedule(accept_connections(server, std::move(*acceptor), false, connections)); + agent_acceptor = std::move(*acceptor); + has_agent_acceptor = true; } else { LOG_WARN("Failed to start agentic listener on {}:{}", opts.host, opts.port); } } - loop.schedule(lsp_peer.run()); + loop.schedule([](MasterServer& server, + kota::ipc::JsonPeer& peer, + std::list& connections, + kota::tcp::acceptor acceptor, + bool has_acceptor) -> kota::task<> { + auto run_peer = [](MasterServer& server, kota::ipc::JsonPeer& peer) -> kota::task<> { + co_await peer.run(); + server.schedule_shutdown(); + }; + auto close_peer_on_shutdown = [](MasterServer& server, + kota::ipc::JsonPeer& peer) -> kota::task<> { + co_await server.get_shutdown_event().wait(); + peer.close(); + }; + + if(has_acceptor) { + co_await kota::when_all( + run_peer(server, peer), + close_peer_on_shutdown(server, peer), + accept_connections(server, std::move(acceptor), false, connections)); + } else { + co_await kota::when_all(run_peer(server, peer), + close_peer_on_shutdown(server, peer)); + } + }(server, lsp_peer, connections, std::move(agent_acceptor), has_agent_acceptor)); loop.run(); return 0; } @@ -485,7 +513,12 @@ static kota::task<> daemon_main(MasterServer& server, kota::pipe::acceptor accep kota::task_group<> connection_group(loop); co_await kota::when_all( - [&]() -> kota::task<> { + [](MasterServer& server, + kota::pipe::acceptor& acceptor, + std::list& connections, + kota::task_group<>& connection_group) -> kota::task<> { + auto& loop = kota::event_loop::current(); + while(true) { auto conn = co_await acceptor.accept(); if(!conn.has_value()) @@ -506,14 +539,16 @@ static kota::task<> daemon_main(MasterServer& server, kota::pipe::acceptor accep connection_group.spawn(run_daemon_connection(peer_ptr, connections, it)); } - }(), - [&]() -> kota::task<> { + }(server, acceptor, connections, connection_group), + [](MasterServer& server, + kota::pipe::acceptor& acceptor, + std::list& connections) -> kota::task<> { co_await server.get_shutdown_event().wait(); acceptor.stop(); for(auto& conn: connections) { conn.peer->close(); } - }()); + }(server, acceptor, connections)); co_await connection_group.join(); } From 7eab896f51126be4c97f9fad9ff8d71a9034734f Mon Sep 17 00:00:00 2001 From: ykiko Date: Fri, 5 Jun 2026 21:46:51 +0800 Subject: [PATCH 03/10] fix(server): simplify shutdown with when_any cancellation propagation Replace explicit peer.close()/acceptor.stop() shutdown handlers with when_any-based cancellation propagation, fix indexer monitor_resources task tracking with dedicated task_group, and bump kotatsu. Co-Authored-By: Claude Opus 4.6 --- cmake/package.cmake | 2 +- src/server/compiler/indexer.cpp | 4 +- src/server/service/agent_client.cpp | 1 - src/server/service/lsp_client.cpp | 1 - src/server/service/master_server.cpp | 279 +++++++++++++-------------- src/server/service/master_server.h | 3 +- 6 files changed, 137 insertions(+), 153 deletions(-) diff --git a/cmake/package.cmake b/cmake/package.cmake index f3ef7b958..5738ee2f1 100644 --- a/cmake/package.cmake +++ b/cmake/package.cmake @@ -41,7 +41,7 @@ set(FLATBUFFERS_BUILD_FLATHASH OFF CACHE BOOL "" FORCE) FetchContent_Declare( kotatsu GIT_REPOSITORY https://github.com/clice-io/kotatsu - GIT_TAG 73814044ce8142f4438a3028f44668675fc09fff + GIT_TAG ac67f39fd3dd2a9daf142c471546980d683e8566 ) set(KOTA_ENABLE_ZEST ON) diff --git a/src/server/compiler/indexer.cpp b/src/server/compiler/indexer.cpp index 4248574ef..dc3cab10e 100644 --- a/src/server/compiler/indexer.cpp +++ b/src/server/compiler/indexer.cpp @@ -889,7 +889,8 @@ kota::task<> Indexer::run_background_indexing() { indexing_active = true; kota::cancellation_source monitor_cancel; - bg_tasks.spawn(kota::with_token(monitor_resources(), monitor_cancel.token())); + kota::task_group<> monitor_group(loop); + monitor_group.spawn(kota::with_token(monitor_resources(), monitor_cancel.token())); std::stable_partition( index_queue.begin() + index_queue_pos, @@ -952,6 +953,7 @@ kota::task<> Indexer::run_background_indexing() { } monitor_cancel.cancel(); + co_await monitor_group.join(); indexing_active = false; LOG_INFO("Background indexing complete: {} files dispatched", dispatched); diff --git a/src/server/service/agent_client.cpp b/src/server/service/agent_client.cpp index d51b610e1..c8eed1d84 100644 --- a/src/server/service/agent_client.cpp +++ b/src/server/service/agent_client.cpp @@ -781,7 +781,6 @@ AgentClient::AgentClient(MasterServer& server, kota::ipc::JsonPeer& peer) : peer.on_notification([this, &srv](const ShutdownParams&) { LOG_INFO("agentic/shutdown received, shutting down"); srv.schedule_shutdown(); - this->peer.close(); }); } diff --git a/src/server/service/lsp_client.cpp b/src/server/service/lsp_client.cpp index b7bbc369c..1607b60a7 100644 --- a/src/server/service/lsp_client.cpp +++ b/src/server/service/lsp_client.cpp @@ -156,7 +156,6 @@ LSPClient::LSPClient(MasterServer& server, kota::ipc::JsonPeer& peer) : server(s peer.on_notification([this]([[maybe_unused]] const protocol::ExitParams& params) { LOG_INFO("Exit notification received"); this->server.schedule_shutdown(); - this->peer.close(); }); peer.on_notification([this](const protocol::DidOpenTextDocumentParams& params) { diff --git a/src/server/service/master_server.cpp b/src/server/service/master_server.cpp index 7ed6cd653..74d4f3155 100644 --- a/src/server/service/master_server.cpp +++ b/src/server/service/master_server.cpp @@ -106,45 +106,40 @@ void MasterServer::initialize(llvm::StringRef root) { initialize(); } -void MasterServer::start_file_watcher() { - if(workspace_root.empty()) - return; - - loop.schedule([](MasterServer& server) -> kota::task<> { - auto watcher = kota::fs_event::create(server.workspace_root, {}, server.loop); - if(!watcher) { - LOG_WARN("Failed to start file watcher for {}", server.workspace_root); - co_return; - } +kota::task<> MasterServer::file_watcher_task() { + auto watcher = kota::fs_event::create(workspace_root, {}, loop); + if(!watcher) { + LOG_WARN("Failed to start file watcher for {}", workspace_root); + co_return; + } - LOG_INFO("File watcher started for {}", server.workspace_root); + LOG_INFO("File watcher started for {}", workspace_root); - while(true) { - auto changes = co_await watcher->next(); - if(!changes) - break; + while(true) { + auto changes = co_await watcher->next(); + if(!changes) + break; - for(auto& change: *changes) { - if(change.type != kota::fs_event::effect::modify && - change.type != kota::fs_event::effect::create) - continue; + for(auto& change: *changes) { + if(change.type != kota::fs_event::effect::modify && + change.type != kota::fs_event::effect::create) + continue; - llvm::StringRef file(change.path); - if(file.ends_with("compile_commands.json")) { - LOG_INFO("CDB changed, reloading workspace"); - server.load_workspace(); - continue; - } + llvm::StringRef file(change.path); + if(file.ends_with("compile_commands.json")) { + LOG_INFO("CDB changed, reloading workspace"); + load_workspace(); + continue; + } - if(file.ends_with(".cpp") || file.ends_with(".cc") || file.ends_with(".cxx") || - file.ends_with(".c") || file.ends_with(".h") || file.ends_with(".hpp") || - file.ends_with(".hxx") || file.ends_with(".cppm") || file.ends_with(".ixx")) { - auto path_id = server.workspace.path_pool.intern(file); - server.on_file_saved(path_id); - } + if(file.ends_with(".cpp") || file.ends_with(".cc") || file.ends_with(".cxx") || + file.ends_with(".c") || file.ends_with(".h") || file.ends_with(".hpp") || + file.ends_with(".hxx") || file.ends_with(".cppm") || file.ends_with(".ixx")) { + auto path_id = workspace.path_pool.intern(file); + on_file_saved(path_id); } } - }(*this)); + } } Session* MasterServer::find_session(std::uint32_t path_id) { @@ -207,15 +202,14 @@ void MasterServer::schedule_shutdown() { if(lifecycle == ServerLifecycle::Exited) return; lifecycle = ServerLifecycle::Exited; + loop.post([this]() { shutdown_event.set(); }); +} +kota::task<> MasterServer::shutdown_and_cleanup() { indexer.save(workspace.config.project.index_dir); workspace.save_cache(); - shutdown_event.set(); - - loop.schedule([](MasterServer& server) -> kota::task<> { - co_await kota::when_all(server.indexer.stop(), server.compiler.stop(), server.pool.stop()); - server.loop.stop(); - }(*this)); + co_await kota::when_all(indexer.stop(), compiler.stop()); + co_await pool.stop(); } void MasterServer::load_workspace() { @@ -354,56 +348,47 @@ static kota::task<> accept_connections(MasterServer& server, bool register_lsp, std::list& connections) { auto& loop = kota::event_loop::current(); - kota::task_group<> connection_group(loop); - - co_await kota::when_all( - [](MasterServer& server, - kota::tcp::acceptor& acceptor, - bool register_lsp, - std::list& connections, - kota::task_group<>& connection_group) -> kota::task<> { - auto& loop = kota::event_loop::current(); - bool lsp_registered = false; - - while(true) { - auto conn = co_await acceptor.accept(); - if(!conn.has_value()) - break; - - LOG_INFO("Client connected"); - - auto transport = std::make_unique(std::move(*conn)); - auto peer = std::make_unique(loop, std::move(transport)); - - std::unique_ptr lsp; - if(register_lsp && !lsp_registered) { - lsp = std::make_unique(server, *peer); - lsp_registered = true; - } - auto agent = std::make_unique(server, *peer); + kota::task_group<> group(loop); + bool lsp_registered = false; - auto* peer_ptr = peer.get(); - auto it = connections.emplace(connections.end(), - Connection{ - .peer = std::move(peer), - .lsp_client = std::move(lsp), - .agent_client = std::move(agent), - }); + group.spawn([](MasterServer& server, + kota::tcp::acceptor& acceptor, + bool register_lsp, + std::list& connections, + kota::task_group<>& group, + bool& lsp_registered) -> kota::task<> { + auto& loop = kota::event_loop::current(); - connection_group.spawn(run_connection(peer_ptr, connections, it)); - } - }(server, acceptor, register_lsp, connections, connection_group), - [](MasterServer& server, - kota::tcp::acceptor& acceptor, - std::list& connections) -> kota::task<> { - co_await server.get_shutdown_event().wait(); - acceptor.stop(); - for(auto& conn: connections) { - conn.peer->close(); + while(true) { + auto conn = co_await acceptor.accept(); + if(!conn.has_value()) + break; + + LOG_INFO("Client connected"); + + auto transport = std::make_unique(std::move(*conn)); + auto peer = std::make_unique(loop, std::move(transport)); + + std::unique_ptr lsp; + if(register_lsp && !lsp_registered) { + lsp = std::make_unique(server, *peer); + lsp_registered = true; } - }(server, acceptor, connections)); + auto agent = std::make_unique(server, *peer); - co_await connection_group.join(); + auto* peer_ptr = peer.get(); + auto it = connections.emplace(connections.end(), + Connection{ + .peer = std::move(peer), + .lsp_client = std::move(lsp), + .agent_client = std::move(agent), + }); + + group.spawn(run_connection(peer_ptr, connections, it)); + } + }(server, acceptor, register_lsp, connections, group, lsp_registered)); + + co_await group.join(); } int run_server_mode(const ServerOptions& opts) { @@ -429,10 +414,6 @@ int run_server_mode(const ServerOptions& opts) { kota::ipc::JsonPeer lsp_peer(loop, std::move(final_transport)); LSPClient lsp_client(server, lsp_peer); - loop.schedule([](MasterServer& server, kota::ipc::JsonPeer& peer) -> kota::task<> { - co_await server.get_shutdown_event().wait(); - peer.close(); - }(server, lsp_peer)); kota::tcp::acceptor agent_acceptor; bool has_agent_acceptor = false; @@ -453,25 +434,15 @@ int run_server_mode(const ServerOptions& opts) { std::list& connections, kota::tcp::acceptor acceptor, bool has_acceptor) -> kota::task<> { - auto run_peer = [](MasterServer& server, kota::ipc::JsonPeer& peer) -> kota::task<> { - co_await peer.run(); - server.schedule_shutdown(); - }; - auto close_peer_on_shutdown = [](MasterServer& server, - kota::ipc::JsonPeer& peer) -> kota::task<> { - co_await server.get_shutdown_event().wait(); - peer.close(); - }; - if(has_acceptor) { - co_await kota::when_all( - run_peer(server, peer), - close_peer_on_shutdown(server, peer), - accept_connections(server, std::move(acceptor), false, connections)); + co_await kota::when_any( + peer.run(), + accept_connections(server, std::move(acceptor), false, connections), + server.get_shutdown_event().wait()); } else { - co_await kota::when_all(run_peer(server, peer), - close_peer_on_shutdown(server, peer)); + co_await kota::when_any(peer.run(), server.get_shutdown_event().wait()); } + co_await server.shutdown_and_cleanup(); }(server, lsp_peer, connections, std::move(agent_acceptor), has_agent_acceptor)); loop.run(); return 0; @@ -485,7 +456,14 @@ int run_server_mode(const ServerOptions& opts) { } LOG_INFO("Listening on {}:{} ...", opts.host, opts.port); - loop.schedule(accept_connections(server, std::move(*acceptor), true, connections)); + loop.schedule([](MasterServer& server, + kota::tcp::acceptor acceptor, + std::list& connections) -> kota::task<> { + co_await kota::when_any( + accept_connections(server, std::move(acceptor), true, connections), + server.get_shutdown_event().wait()); + co_await server.shutdown_and_cleanup(); + }(server, std::move(*acceptor), connections)); loop.run(); return 0; } @@ -507,50 +485,54 @@ static kota::task<> run_daemon_connection(kota::ipc::JsonPeer* peer, connections.erase(pos); } -static kota::task<> daemon_main(MasterServer& server, kota::pipe::acceptor acceptor) { +static kota::task<> daemon_accept(MasterServer& server, kota::pipe::acceptor acceptor) { auto& loop = kota::event_loop::current(); std::list connections; - kota::task_group<> connection_group(loop); - - co_await kota::when_all( - [](MasterServer& server, - kota::pipe::acceptor& acceptor, - std::list& connections, - kota::task_group<>& connection_group) -> kota::task<> { - auto& loop = kota::event_loop::current(); - - while(true) { - auto conn = co_await acceptor.accept(); - if(!conn.has_value()) - break; - - LOG_INFO("Daemon client connected"); - - auto transport = std::make_unique(std::move(*conn)); - auto peer = std::make_unique(loop, std::move(transport)); - auto agent = std::make_unique(server, *peer); - - auto* peer_ptr = peer.get(); - auto it = connections.emplace(connections.end(), - DaemonConnection{ - .peer = std::move(peer), - .agent_client = std::move(agent), - }); - - connection_group.spawn(run_daemon_connection(peer_ptr, connections, it)); - } - }(server, acceptor, connections, connection_group), - [](MasterServer& server, - kota::pipe::acceptor& acceptor, - std::list& connections) -> kota::task<> { - co_await server.get_shutdown_event().wait(); - acceptor.stop(); - for(auto& conn: connections) { - conn.peer->close(); - } - }(server, acceptor, connections)); + kota::task_group<> group(loop); + + group.spawn([](MasterServer& server, + kota::pipe::acceptor& acceptor, + std::list& connections, + kota::task_group<>& group) -> kota::task<> { + auto& loop = kota::event_loop::current(); + + while(true) { + auto conn = co_await acceptor.accept(); + if(!conn.has_value()) + break; + + LOG_INFO("Daemon client connected"); + + auto transport = std::make_unique(std::move(*conn)); + auto peer = std::make_unique(loop, std::move(transport)); + auto agent = std::make_unique(server, *peer); + + auto* peer_ptr = peer.get(); + auto it = connections.emplace(connections.end(), + DaemonConnection{ + .peer = std::move(peer), + .agent_client = std::move(agent), + }); - co_await connection_group.join(); + group.spawn(run_daemon_connection(peer_ptr, connections, it)); + } + }(server, acceptor, connections, group)); + + co_await group.join(); +} + +static kota::task<> daemon_main(MasterServer& server, + kota::pipe::acceptor acceptor, + bool watch_files) { + if(watch_files) { + co_await kota::when_any(daemon_accept(server, std::move(acceptor)), + server.file_watcher_task(), + server.get_shutdown_event().wait()); + } else { + co_await kota::when_any(daemon_accept(server, std::move(acceptor)), + server.get_shutdown_event().wait()); + } + co_await server.shutdown_and_cleanup(); } int run_daemon_mode(const DaemonOptions& opts) { @@ -586,9 +568,10 @@ int run_daemon_mode(const DaemonOptions& opts) { kota::event_loop loop; MasterServer server(loop, opts.self_path); + bool watch_files = false; if(!opts.workspace.empty()) { server.initialize(opts.workspace); - server.start_file_watcher(); + watch_files = true; } auto acceptor = kota::pipe::listen(socket_path, {}, loop); @@ -598,7 +581,7 @@ int run_daemon_mode(const DaemonOptions& opts) { } LOG_INFO("Daemon listening on {}", socket_path); - loop.schedule(daemon_main(server, std::move(*acceptor))); + loop.schedule(daemon_main(server, std::move(*acceptor), watch_files)); loop.run(); llvm::sys::fs::remove(socket_path); diff --git a/src/server/service/master_server.h b/src/server/service/master_server.h index 2a9f37938..bc76202fa 100644 --- a/src/server/service/master_server.h +++ b/src/server/service/master_server.h @@ -39,7 +39,8 @@ class MasterServer { void initialize(); void initialize(llvm::StringRef root); - void start_file_watcher(); + kota::task<> file_watcher_task(); + kota::task<> shutdown_and_cleanup(); Session* find_session(std::uint32_t path_id); Session& open_session(std::uint32_t path_id); From ca616ee1692e158116b6a37476b09f0d650690d8 Mon Sep 17 00:00:00 2001 From: ykiko Date: Sun, 7 Jun 2026 18:32:49 +0800 Subject: [PATCH 04/10] fix(server): structured shutdown with kotatsu cancel fix Update kotatsu to ea7d99b which fixes cancel propagation for reentrantly-cancelled tasks (#158, #160). Use when_any-based structured shutdown so all tasks are properly cancelled before cleanup, avoiding ASAN use-after-free on exit. Add close_peer_on_shutdown workaround for pipe mode: the sentinel prevents cancel from reaching pending I/O when shutdown fires inline from peer.run()'s exit handler. --- cmake/package.cmake | 2 +- src/server/service/master_server.cpp | 48 ++++++++++++++++++---------- 2 files changed, 32 insertions(+), 18 deletions(-) diff --git a/cmake/package.cmake b/cmake/package.cmake index 5738ee2f1..b52184bfe 100644 --- a/cmake/package.cmake +++ b/cmake/package.cmake @@ -41,7 +41,7 @@ set(FLATBUFFERS_BUILD_FLATHASH OFF CACHE BOOL "" FORCE) FetchContent_Declare( kotatsu GIT_REPOSITORY https://github.com/clice-io/kotatsu - GIT_TAG ac67f39fd3dd2a9daf142c471546980d683e8566 + GIT_TAG ea7d99b2b703c813a6246129934119898fa58da5 ) set(KOTA_ENABLE_ZEST ON) diff --git a/src/server/service/master_server.cpp b/src/server/service/master_server.cpp index 74d4f3155..4df0c6772 100644 --- a/src/server/service/master_server.cpp +++ b/src/server/service/master_server.cpp @@ -202,7 +202,7 @@ void MasterServer::schedule_shutdown() { if(lifecycle == ServerLifecycle::Exited) return; lifecycle = ServerLifecycle::Exited; - loop.post([this]() { shutdown_event.set(); }); + shutdown_event.set(); } kota::task<> MasterServer::shutdown_and_cleanup() { @@ -335,12 +335,9 @@ struct Connection { std::unique_ptr agent_client; }; -static kota::task<> run_connection(kota::ipc::JsonPeer* peer, - std::list& connections, - std::list::iterator pos) { +static kota::task<> run_connection(kota::ipc::JsonPeer* peer) { co_await peer->run(); LOG_INFO("Client disconnected"); - connections.erase(pos); } static kota::task<> accept_connections(MasterServer& server, @@ -384,7 +381,7 @@ static kota::task<> accept_connections(MasterServer& server, .agent_client = std::move(agent), }); - group.spawn(run_connection(peer_ptr, connections, it)); + group.spawn(run_connection(peer_ptr)); } }(server, acceptor, register_lsp, connections, group, lsp_registered)); @@ -429,21 +426,42 @@ int run_server_mode(const ServerOptions& opts) { } } + // FIXME(kotatsu#160): workaround for incomplete cancel propagation. + // shutdown_event.set() fires inline from within peer.run() (exit + // notification handler), so when_any's cancel hits the sentinel + // (child==self) and never reaches the pending transport read. + // The task's state IS set to Cancelled, but the io_op is unaware + // and blocks forever. Explicitly closing the peer forces the read + // to fail, allowing on_child_complete's Cancelled check to finalize. + // Ideally when_any would cancel I/O children of a sentinel-guarded + // task once the sentinel is cleared (i.e. at the next co_await). + auto close_peer_on_shutdown = [](MasterServer& server, + kota::ipc::JsonPeer& peer) -> kota::task<> { + co_await server.get_shutdown_event().wait(); + peer.close(); + }; + loop.schedule([](MasterServer& server, kota::ipc::JsonPeer& peer, std::list& connections, kota::tcp::acceptor acceptor, - bool has_acceptor) -> kota::task<> { + bool has_acceptor, + auto close_peer_on_shutdown) -> kota::task<> { if(has_acceptor) { co_await kota::when_any( peer.run(), - accept_connections(server, std::move(acceptor), false, connections), - server.get_shutdown_event().wait()); + close_peer_on_shutdown(server, peer), + accept_connections(server, std::move(acceptor), false, connections)); } else { - co_await kota::when_any(peer.run(), server.get_shutdown_event().wait()); + co_await kota::when_any(peer.run(), close_peer_on_shutdown(server, peer)); } co_await server.shutdown_and_cleanup(); - }(server, lsp_peer, connections, std::move(agent_acceptor), has_agent_acceptor)); + }(server, + lsp_peer, + connections, + std::move(agent_acceptor), + has_agent_acceptor, + close_peer_on_shutdown)); loop.run(); return 0; } @@ -477,19 +495,15 @@ struct DaemonConnection { std::unique_ptr agent_client; }; -static kota::task<> run_daemon_connection(kota::ipc::JsonPeer* peer, - std::list& connections, - std::list::iterator pos) { +static kota::task<> run_daemon_connection(kota::ipc::JsonPeer* peer) { co_await peer->run(); LOG_INFO("Daemon client disconnected"); - connections.erase(pos); } static kota::task<> daemon_accept(MasterServer& server, kota::pipe::acceptor acceptor) { auto& loop = kota::event_loop::current(); std::list connections; kota::task_group<> group(loop); - group.spawn([](MasterServer& server, kota::pipe::acceptor& acceptor, std::list& connections, @@ -514,7 +528,7 @@ static kota::task<> daemon_accept(MasterServer& server, kota::pipe::acceptor acc .agent_client = std::move(agent), }); - group.spawn(run_daemon_connection(peer_ptr, connections, it)); + group.spawn(run_daemon_connection(peer_ptr)); } }(server, acceptor, connections, group)); From 52361e3f8d769959af8ce40991ce976938aeac54 Mon Sep 17 00:00:00 2001 From: ykiko Date: Sun, 7 Jun 2026 21:47:36 +0800 Subject: [PATCH 05/10] fix(server): remove shutdown workaround with kotatsu deferred dispatch kotatsu 5e059a2 defers sync primitive resumes to the event loop idle tick, eliminating inline reentrancy. cancel() now propagates cleanly through when_any to the transport io_op, so the close_peer_on_shutdown workaround is no longer needed. --- cmake/package.cmake | 2 +- src/server/service/master_server.cpp | 31 +++++----------------------- 2 files changed, 6 insertions(+), 27 deletions(-) diff --git a/cmake/package.cmake b/cmake/package.cmake index b52184bfe..c49bcbe78 100644 --- a/cmake/package.cmake +++ b/cmake/package.cmake @@ -41,7 +41,7 @@ set(FLATBUFFERS_BUILD_FLATHASH OFF CACHE BOOL "" FORCE) FetchContent_Declare( kotatsu GIT_REPOSITORY https://github.com/clice-io/kotatsu - GIT_TAG ea7d99b2b703c813a6246129934119898fa58da5 + GIT_TAG 5e059a2d346200212f465094babde0b40b71a6f6 ) set(KOTA_ENABLE_ZEST ON) diff --git a/src/server/service/master_server.cpp b/src/server/service/master_server.cpp index 4df0c6772..73ab6b9b9 100644 --- a/src/server/service/master_server.cpp +++ b/src/server/service/master_server.cpp @@ -426,42 +426,21 @@ int run_server_mode(const ServerOptions& opts) { } } - // FIXME(kotatsu#160): workaround for incomplete cancel propagation. - // shutdown_event.set() fires inline from within peer.run() (exit - // notification handler), so when_any's cancel hits the sentinel - // (child==self) and never reaches the pending transport read. - // The task's state IS set to Cancelled, but the io_op is unaware - // and blocks forever. Explicitly closing the peer forces the read - // to fail, allowing on_child_complete's Cancelled check to finalize. - // Ideally when_any would cancel I/O children of a sentinel-guarded - // task once the sentinel is cleared (i.e. at the next co_await). - auto close_peer_on_shutdown = [](MasterServer& server, - kota::ipc::JsonPeer& peer) -> kota::task<> { - co_await server.get_shutdown_event().wait(); - peer.close(); - }; - loop.schedule([](MasterServer& server, kota::ipc::JsonPeer& peer, std::list& connections, kota::tcp::acceptor acceptor, - bool has_acceptor, - auto close_peer_on_shutdown) -> kota::task<> { + bool has_acceptor) -> kota::task<> { if(has_acceptor) { co_await kota::when_any( peer.run(), - close_peer_on_shutdown(server, peer), - accept_connections(server, std::move(acceptor), false, connections)); + accept_connections(server, std::move(acceptor), false, connections), + server.get_shutdown_event().wait()); } else { - co_await kota::when_any(peer.run(), close_peer_on_shutdown(server, peer)); + co_await kota::when_any(peer.run(), server.get_shutdown_event().wait()); } co_await server.shutdown_and_cleanup(); - }(server, - lsp_peer, - connections, - std::move(agent_acceptor), - has_agent_acceptor, - close_peer_on_shutdown)); + }(server, lsp_peer, connections, std::move(agent_acceptor), has_agent_acceptor)); loop.run(); return 0; } From 4a3428e0658c43a41406d5e29bd50aceb110bed2 Mon Sep 17 00:00:00 2001 From: ykiko Date: Sun, 7 Jun 2026 22:45:27 +0800 Subject: [PATCH 06/10] fix(server): address review issues from shutdown refactoring - Fix lifecycle state: schedule_shutdown sets ShuttingDown (not Exited), shutdown_and_cleanup sets Exited after cleanup completes. Guard only checks Exited to avoid blocking the exit notification when the LSP shutdown request already set ShuttingDown. - Restore connection cleanup: erase Connection from list when peer disconnects to prevent unbounded accumulation. - Wrap file_watcher_task in resilient_file_watcher so watcher creation failure doesn't become a when_any winner that shuts down the daemon. - Remove unused `this` capture in agent_client shutdown handler. --- src/server/service/agent_client.cpp | 2 +- src/server/service/master_server.cpp | 24 ++++++++++++++++++------ 2 files changed, 19 insertions(+), 7 deletions(-) diff --git a/src/server/service/agent_client.cpp b/src/server/service/agent_client.cpp index c8eed1d84..69c7a81b2 100644 --- a/src/server/service/agent_client.cpp +++ b/src/server/service/agent_client.cpp @@ -778,7 +778,7 @@ AgentClient::AgentClient(MasterServer& server, kota::ipc::JsonPeer& peer) : co_return result; }); - peer.on_notification([this, &srv](const ShutdownParams&) { + peer.on_notification([&srv](const ShutdownParams&) { LOG_INFO("agentic/shutdown received, shutting down"); srv.schedule_shutdown(); }); diff --git a/src/server/service/master_server.cpp b/src/server/service/master_server.cpp index 73ab6b9b9..bf22a36c4 100644 --- a/src/server/service/master_server.cpp +++ b/src/server/service/master_server.cpp @@ -201,7 +201,7 @@ void MasterServer::on_file_saved(std::uint32_t path_id) { void MasterServer::schedule_shutdown() { if(lifecycle == ServerLifecycle::Exited) return; - lifecycle = ServerLifecycle::Exited; + lifecycle = ServerLifecycle::ShuttingDown; shutdown_event.set(); } @@ -210,6 +210,7 @@ kota::task<> MasterServer::shutdown_and_cleanup() { workspace.save_cache(); co_await kota::when_all(indexer.stop(), compiler.stop()); co_await pool.stop(); + lifecycle = ServerLifecycle::Exited; } void MasterServer::load_workspace() { @@ -335,9 +336,12 @@ struct Connection { std::unique_ptr agent_client; }; -static kota::task<> run_connection(kota::ipc::JsonPeer* peer) { +static kota::task<> run_connection(kota::ipc::JsonPeer* peer, + std::list& connections, + std::list::iterator pos) { co_await peer->run(); LOG_INFO("Client disconnected"); + connections.erase(pos); } static kota::task<> accept_connections(MasterServer& server, @@ -381,7 +385,7 @@ static kota::task<> accept_connections(MasterServer& server, .agent_client = std::move(agent), }); - group.spawn(run_connection(peer_ptr)); + group.spawn(run_connection(peer_ptr, connections, it)); } }(server, acceptor, register_lsp, connections, group, lsp_registered)); @@ -474,9 +478,12 @@ struct DaemonConnection { std::unique_ptr agent_client; }; -static kota::task<> run_daemon_connection(kota::ipc::JsonPeer* peer) { +static kota::task<> run_daemon_connection(kota::ipc::JsonPeer* peer, + std::list& connections, + std::list::iterator pos) { co_await peer->run(); LOG_INFO("Daemon client disconnected"); + connections.erase(pos); } static kota::task<> daemon_accept(MasterServer& server, kota::pipe::acceptor acceptor) { @@ -507,19 +514,24 @@ static kota::task<> daemon_accept(MasterServer& server, kota::pipe::acceptor acc .agent_client = std::move(agent), }); - group.spawn(run_daemon_connection(peer_ptr)); + group.spawn(run_daemon_connection(peer_ptr, connections, it)); } }(server, acceptor, connections, group)); co_await group.join(); } +static kota::task<> resilient_file_watcher(MasterServer& server) { + co_await server.file_watcher_task(); + co_await server.get_shutdown_event().wait(); +} + static kota::task<> daemon_main(MasterServer& server, kota::pipe::acceptor acceptor, bool watch_files) { if(watch_files) { co_await kota::when_any(daemon_accept(server, std::move(acceptor)), - server.file_watcher_task(), + resilient_file_watcher(server), server.get_shutdown_event().wait()); } else { co_await kota::when_any(daemon_accept(server, std::move(acceptor)), From 549f208adf38ebe2bb9c3a4af2e1c3faeeb7f897 Mon Sep 17 00:00:00 2001 From: ykiko Date: Mon, 8 Jun 2026 12:23:47 +0800 Subject: [PATCH 07/10] chore: bump kotatsu to eabd6c1 (deferred resume improvements) Picks up stabilized cancellation handling, grant abandonment for mutex/semaphore, and immediate drain of deferred resumes after the outermost coroutine resume returns. --- cmake/package.cmake | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cmake/package.cmake b/cmake/package.cmake index c49bcbe78..b3bc885e8 100644 --- a/cmake/package.cmake +++ b/cmake/package.cmake @@ -41,7 +41,7 @@ set(FLATBUFFERS_BUILD_FLATHASH OFF CACHE BOOL "" FORCE) FetchContent_Declare( kotatsu GIT_REPOSITORY https://github.com/clice-io/kotatsu - GIT_TAG 5e059a2d346200212f465094babde0b40b71a6f6 + GIT_TAG 2a8c147579d36d33ff50f31d0bce49b0ab67381b ) set(KOTA_ENABLE_ZEST ON) From 55eed02d6d43e48f1f77843477013864f5072de5 Mon Sep 17 00:00:00 2001 From: ykiko Date: Mon, 8 Jun 2026 23:14:05 +0800 Subject: [PATCH 08/10] fix(server): address review issues from shutdown refactoring Revert indexer monitor task to bg_tasks.spawn() to prevent UAF when background indexing is cancelled, and improve test_rpc_shutdown to check for sanitizer errors via assert_server_exited_cleanly(). --- src/server/compiler/indexer.cpp | 4 +--- tests/integration/agentic/test_agentic.py | 15 ++++++--------- 2 files changed, 7 insertions(+), 12 deletions(-) diff --git a/src/server/compiler/indexer.cpp b/src/server/compiler/indexer.cpp index dc3cab10e..4248574ef 100644 --- a/src/server/compiler/indexer.cpp +++ b/src/server/compiler/indexer.cpp @@ -889,8 +889,7 @@ kota::task<> Indexer::run_background_indexing() { indexing_active = true; kota::cancellation_source monitor_cancel; - kota::task_group<> monitor_group(loop); - monitor_group.spawn(kota::with_token(monitor_resources(), monitor_cancel.token())); + bg_tasks.spawn(kota::with_token(monitor_resources(), monitor_cancel.token())); std::stable_partition( index_queue.begin() + index_queue_pos, @@ -953,7 +952,6 @@ kota::task<> Indexer::run_background_indexing() { } monitor_cancel.cancel(); - co_await monitor_group.join(); indexing_active = false; LOG_INFO("Background indexing complete: {} files dispatched", dispatched); diff --git a/tests/integration/agentic/test_agentic.py b/tests/integration/agentic/test_agentic.py index 75ad08157..5d529fd85 100644 --- a/tests/integration/agentic/test_agentic.py +++ b/tests/integration/agentic/test_agentic.py @@ -422,9 +422,9 @@ async def test_rpc_status(indexed_agentic, workspace): @pytest.mark.workspace("hello_world") async def test_rpc_shutdown(executable, workspace): - """Shutdown notification should cause the server to exit.""" + """Shutdown notification should cause the server to exit cleanly.""" from tests.integration.utils.client import CliceClient - from tests.conftest import _shutdown_client, _find_free_port + from tests.conftest import _find_free_port, assert_server_exited_cleanly host = "127.0.0.1" port = _find_free_port() @@ -445,13 +445,10 @@ async def test_rpc_shutdown(executable, workspace): pass rpc.sock.close() - import asyncio - - for _ in range(20): - if c._server.returncode is not None: - break - await asyncio.sleep(0.5) - assert c._server.returncode is not None, "Server did not exit after shutdown" + await assert_server_exited_cleanly(c._server) + c._stop_event.set() + for task in c._async_tasks: + task.cancel() @pytest.mark.workspace("index_features") From cf7d77475d5c705dc0b096b3ec6781bdb5e80fa1 Mon Sep 17 00:00:00 2001 From: ykiko Date: Mon, 8 Jun 2026 23:40:08 +0800 Subject: [PATCH 09/10] fix(tests): increase shutdown timeout for slow CI runners The 3s default was too tight for CI environments, especially Debug builds where indexer/compiler cleanup takes longer. Bumped to 10s. --- tests/conftest.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/conftest.py b/tests/conftest.py index ecca7a0ad..8e77755b2 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -209,7 +209,7 @@ def _server_stderr_excerpt(stderr_text: str) -> str: return "\n".join(interesting[-80:]) -async def assert_server_exited_cleanly(server, timeout: float = 3.0) -> None: +async def assert_server_exited_cleanly(server, timeout: float = 10.0) -> None: failures: list[str] = [] if server is None: @@ -254,7 +254,7 @@ async def _shutdown_client(c: CliceClient) -> None: server = getattr(c, "_server", None) try: - await asyncio.wait_for(c.shutdown_async(None), timeout=3.0) + await asyncio.wait_for(c.shutdown_async(None), timeout=10.0) except Exception: pass From 4fc239f4130ad29f0c8dcfea341dde18a36db52d Mon Sep 17 00:00:00 2001 From: ykiko Date: Mon, 8 Jun 2026 23:53:46 +0800 Subject: [PATCH 10/10] fix(tests): suppress pygls DeprecationWarning in pytest output --- tests/pytest.ini | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tests/pytest.ini b/tests/pytest.ini index db904db7c..ebf65ab17 100644 --- a/tests/pytest.ini +++ b/tests/pytest.ini @@ -1,5 +1,7 @@ [pytest] asyncio_mode = auto +filterwarnings = + ignore::DeprecationWarning:pygls markers = workspace init_options