Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
2 changes: 1 addition & 1 deletion cmake/package.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -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 2a8c147579d36d33ff50f31d0bce49b0ab67381b
)

set(KOTA_ENABLE_ZEST ON)
Expand Down
239 changes: 142 additions & 97 deletions src/server/service/master_server.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -106,45 +106,40 @@ void MasterServer::initialize(llvm::StringRef root) {
initialize();
}

void MasterServer::start_file_watcher() {
if(workspace_root.empty())
return;

loop.schedule([this]() -> kota::task<> {
auto watcher = kota::fs_event::create(workspace_root, {}, loop);
if(!watcher) {
LOG_WARN("Failed to start file watcher for {}", 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 {}", 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");
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 = workspace.path_pool.intern(file);
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);
}
}
}());
}
}

Session* MasterServer::find_session(std::uint32_t path_id) {
Expand Down Expand Up @@ -206,16 +201,16 @@ 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();
}

kota::task<> MasterServer::shutdown_and_cleanup() {
indexer.save(workspace.config.project.index_dir);
workspace.save_cache();
shutdown_event.set();

loop.schedule([this]() -> kota::task<> {
co_await kota::when_all(indexer.stop(), compiler.stop(), pool.stop());
loop.stop();
}());
co_await kota::when_all(indexer.stop(), compiler.stop());
co_await pool.stop();
lifecycle = ServerLifecycle::Exited;
}

void MasterServer::load_workspace() {
Expand Down Expand Up @@ -354,38 +349,47 @@ static kota::task<> accept_connections(MasterServer& server,
bool register_lsp,
std::list<Connection>& connections) {
auto& loop = kota::event_loop::current();
kota::task_group<> connection_group(loop);
kota::task_group<> group(loop);
bool lsp_registered = false;

while(true) {
auto conn = co_await acceptor.accept();
if(!conn.has_value())
break;
group.spawn([](MasterServer& server,
kota::tcp::acceptor& acceptor,
bool register_lsp,
std::list<Connection>& connections,
kota::task_group<>& group,
bool& lsp_registered) -> kota::task<> {
auto& loop = kota::event_loop::current();

LOG_INFO("Client connected");
while(true) {
auto conn = co_await acceptor.accept();
if(!conn.has_value())
break;

auto transport = std::make_unique<kota::ipc::StreamTransport>(std::move(*conn));
auto peer = std::make_unique<kota::ipc::JsonPeer>(loop, std::move(transport));
LOG_INFO("Client connected");

std::unique_ptr<LSPClient> lsp;
if(register_lsp && !lsp_registered) {
lsp = std::make_unique<LSPClient>(server, *peer);
lsp_registered = true;
}
auto agent = std::make_unique<AgentClient>(server, *peer);
auto transport = std::make_unique<kota::ipc::StreamTransport>(std::move(*conn));
auto peer = std::make_unique<kota::ipc::JsonPeer>(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<LSPClient> lsp;
if(register_lsp && !lsp_registered) {
lsp = std::make_unique<LSPClient>(server, *peer);
lsp_registered = true;
}
auto agent = std::make_unique<AgentClient>(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),
});

co_await connection_group.join();
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) {
Expand All @@ -412,17 +416,35 @@ int run_server_mode(const ServerOptions& opts) {
kota::ipc::JsonPeer lsp_peer(loop, std::move(final_transport));
LSPClient lsp_client(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<Connection>& connections,
kota::tcp::acceptor acceptor,
bool has_acceptor) -> 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());
} else {
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;
}
Expand All @@ -435,7 +457,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<Connection>& 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;
}
Expand All @@ -457,43 +486,58 @@ 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<DaemonConnection> connections;
kota::task_group<> connection_group(loop);
kota::task_group<> group(loop);
group.spawn([](MasterServer& server,
kota::pipe::acceptor& acceptor,
std::list<DaemonConnection>& connections,
kota::task_group<>& group) -> kota::task<> {
auto& loop = kota::event_loop::current();

co_await kota::when_all(
[&]() -> kota::task<> {
while(true) {
auto conn = co_await acceptor.accept();
if(!conn.has_value())
break;
while(true) {
auto conn = co_await acceptor.accept();
if(!conn.has_value())
break;

LOG_INFO("Daemon client connected");
LOG_INFO("Daemon client connected");

auto transport = std::make_unique<kota::ipc::StreamTransport>(std::move(*conn));
auto peer = std::make_unique<kota::ipc::JsonPeer>(loop, std::move(transport));
auto agent = std::make_unique<AgentClient>(server, *peer);
auto transport = std::make_unique<kota::ipc::StreamTransport>(std::move(*conn));
auto peer = std::make_unique<kota::ipc::JsonPeer>(loop, std::move(transport));
auto agent = std::make_unique<AgentClient>(server, *peer);

auto* peer_ptr = peer.get();
auto it = connections.emplace(connections.end(),
DaemonConnection{
.peer = std::move(peer),
.agent_client = std::move(agent),
});
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));
}
}(),
[&]() -> kota::task<> {
co_await server.get_shutdown_event().wait();
acceptor.stop();
for(auto& conn: connections) {
conn.peer->close();
}
}());
group.spawn(run_daemon_connection(peer_ptr, connections, it));
}
}(server, acceptor, connections, group));

co_await connection_group.join();
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)),
resilient_file_watcher(server),
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) {
Expand Down Expand Up @@ -529,9 +573,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);
Expand All @@ -541,7 +586,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);
Expand Down
3 changes: 2 additions & 1 deletion src/server/service/master_server.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
7 changes: 5 additions & 2 deletions src/server/worker/worker_pool.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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_)
Expand Down
Loading
Loading