//! A real Language Server Protocol client: child processes spoken to over //! JSON-RPC 2.0 with `Content-Length` framing. Nothing here knows any single //! language — `specs` is a table of (binary, languageId, extensions, root //! markers), and rust-analyzer, clangd and gopls are rows in it. The in-process //! ZLS backend keeps `.zig`; this file is every language pardes highlights but //! could not answer questions about. //! //! THE PROCESS LIFECYCLE IS THE DESIGN. `lsp.query` is a synchronous call on a //! worker thread, and a language server costs tens of milliseconds to start //! and MINUTES to index a large workspace. So a server is spawned lazily on //! the first query that needs it, kept for the life of the editor, and each //! spec gets at most one — `conns[i]` guards itself with a pthread mutex //! because the gui shell detaches its workers and a superseded query can still //! be inside `run` when the next arrives. //! //! ONE READER THREAD PER SERVER, and it is not optional. The first design //! pumped the socket only while a query waited, which works for a server that //! only ever answers. A real server TALKS: rust-analyzer streams `$/progress` //! for the whole minutes-long index of a big workspace, publishes diagnostics //! it was never asked for, and asks its own `workspace/configuration` //! questions mid-flight. The reader owns the read side of the socket, routes //! responses to the one waiting query (a mailbox under the conn's mutex), //! answers server-to-client requests so the server never blocks on us, feeds //! the diagnostics store, and narrates state changes through `status sink` — //! the shell posts them to `Pardes.setMessage`, so "rust-analyzer indexing 45%" //! lands on the same transient message row a save narrates into. The reader is //! also the ONLY closer of its socket fd: teardown calls `shutdown(2)` and the //! reader closes on the EOF it then reads, so the fd number cannot be recycled //! under a thread still polling it. //! //! Nothing here may wedge the editor: //! - every write and every mailbox wait is bounded by a deadline, //! - a query the server does not answer in time returns no rows and sends //! `$/cancelRequest`; the server is NOT killed for being busy (an indexing //! server is busy for minutes and the status row says so) — but three //! consecutive timeouts mean wedged, and wedged is killed and respawned, //! - the transport is an AF_UNIX socketpair, so a dead server answers EPIPE //! from `send(MSG_NOSIGNAL)` instead of killing pardes with SIGPIPE //! (ignoring SIGPIPE process-wide would be inherited by every pty shell we //! fork and would change what `yes | head` does in a pane), //! - a binary missing from PATH costs one cached probe, disables the spec //! for the session, and says so ONCE on the message row. //! //! Position encoding is negotiated to utf-8 and the SERVER'S ANSWER is //! believed, not our request; the utf-16 conversion is implemented in both //! directions for servers that refuse (`Req.offset` is a byte offset, and a //! misconverted column is wrong on every line with non-ASCII in it). const std = @import("std"); const builtin = @import("builtin"); const libc = std.c; const lsp = @import("lsp.zig"); extern "c" fn execv(path: [*:0]const u8, argv: [*:null]const ?[*:0]const u8) c_int; extern "c" fn chdir(path: [*:0]const u8) c_int; extern "c" fn _exit(status: c_int) noreturn; extern "c" fn setsid() libc.pid_t; extern "c" fn access(path: [*:0]const u8, mode: c_int) c_int; /// std.posix.getenv is gone in 0.16 and std.process.Environ wants an Io extern "c" fn getenv(name: [*:0]const u8) ?[*:0]const u8; extern "c" fn usleep(usec: c_uint) c_int; extern "c" fn realpath(path: [*:0]const u8, resolved: [*]u8) ?[*:0]u8; // test-only libc (the seam's own tests build a real directory tree) extern "c" fn mkdtemp(template: [*:0]u8) ?[*:0]u8; extern "c" fn mkdir(path: [*:0]const u8, mode: libc.mode_t) c_int; extern "c" fn system(cmd: [*:0]const u8) c_int; // pthread condattr surface, absent from std.c: what makes the mailbox waits // tick on CLOCK_MONOTONIC (linux) or a relative timeout (darwin). const pthread_condattr_t = extern struct { data: [8]u8 align(@alignOf(usize)) = @splat(0) }; extern "c" fn pthread_condattr_init(attr: *pthread_condattr_t) c_int; extern "c" fn pthread_condattr_setclock(attr: *pthread_condattr_t, clock: c_int) c_int; extern "c" fn pthread_condattr_destroy(attr: *pthread_condattr_t) c_int; extern "c" fn pthread_cond_init(cond: *libc.pthread_cond_t, attr: ?*const pthread_condattr_t) c_int; extern "c" fn pthread_cond_timedwait_relative_np(cond: *libc.pthread_cond_t, mutex: *libc.pthread_mutex_t, reltime: *const libc.timespec) c_int; const CLOCK_MONOTONIC: c_int = 1; // linux ABI; the setclock call is linux-only const X_OK: c_int = 1; const POLLIN: i16 = 0x001; const POLLOUT: i16 = 0x004; /// linux MSG_NOSIGNAL; darwin has no such flag and gets SO_NOSIGPIPE on the /// socket instead (0x1022), set right after socketpair. Both numbers are ABI. const msg_nosignal: c_int = if (builtin.os.tag.isDarwin()) 0 else 0x4000; const so_nosigpipe: c_int = 0x1022; /// One language server this client knows how to run. Adding a language is /// adding a row; nothing below the table branches on a language. const Spec = struct { /// what the message row and `SPC l i` call it name: []const u8, /// argv[0], searched on PATH unless it contains a slash bin: []const u8, args: []const []const u8 = &.{}, /// the protocol's `languageId` for didOpen lang: []const u8, /// extensions that route a file here (the seam asks `speaks`) exts: []const []const u8, /// project markers, walked UP from the file: the TOP-MOST directory /// holding one is the root (helix's find_root rule), `.git` the fallback markers: []const []const u8, /// environment override: a binary path/name to use instead of `bin`, or /// empty ("") to disable the spec entirely. How the snapshot harness pins /// a deterministic mock server, and how a user points at a custom build. env: [:0]const u8, }; pub const specs = [_]Spec{ .{ .name = "rust-analyzer", .bin = "rust-analyzer", .lang = "rust", .exts = &.{".rs"}, .markers = &.{ "Cargo.toml", "rust-project.json" }, .env = "PARDES_LSP_RS", }, .{ .name = "clangd", .bin = "clangd", .lang = "c", .exts = &.{ ".c", ".h", ".cc", ".cpp", ".hpp", ".cxx", ".hxx" }, .markers = &.{ "compile_commands.json", "compile_flags.txt", ".clangd" }, .env = "PARDES_LSP_C", }, .{ .name = "gopls", .bin = "gopls", .lang = "go", .exts = &.{".go"}, .markers = &.{ "go.mod", "go.work" }, .env = "PARDES_LSP_GO", }, .{ .name = "typescript-language-server", .bin = "typescript-language-server", .args = &.{"--stdio"}, .lang = "typescript", .exts = &.{ ".ts", ".tsx", ".js", ".jsx", ".mjs", ".cjs" }, .markers = &.{ "tsconfig.json", "jsconfig.json", "package.json" }, .env = "PARDES_LSP_TS", }, .{ .name = "pyright", .bin = "pyright-langserver", .args = &.{"--stdio"}, .lang = "python", .exts = &.{".py"}, .markers = &.{ "pyproject.toml", "setup.py", "requirements.txt" }, .env = "PARDES_LSP_PY", }, }; /// Everything, including the four hierarchy kinds the in-process backend has /// no analyser for. Whether one SERVER can answer is a capability question /// answered per connection; a kind its server never advertised simply returns /// no rows, which the harness reports as CLAIMED-EMPTY per language — honest, /// since the claim here is about the protocol, not about every server. pub const supports: std.EnumSet(lsp.Kind) = .initFull(); /// Routing: an extension in the table whose spec is not disabled by its env /// var. Deliberately does NOT probe for the binary — this runs on the Tab /// keystroke. A missing binary is discovered at spawn, disables the spec, and /// says so once on the message row; until then Tab in a `.rs` file diverts, /// gets "no rows" instantly (the disabled flag short-circuits), and indents /// late exactly like any other unanswerable completion. pub fn speaks(path: []const u8) bool { return specFor(path) != null; } fn specFor(path: []const u8) ?usize { for (&specs, 0..) |*s, i| { for (s.exts) |e| { if (!std.mem.endsWith(u8, path, e)) continue; if (getenv(s.env)) |v| if (v[0] == 0) return null; // "" disables return i; } } return null; } // Deadlines. Requests are bounded because tty.zig cancels-and-joins the // previous worker on a new keypress, so the worst UI stall a wedged wait can // cause is one req_ms. Timeouts do NOT kill the server — an indexing // rust-analyzer legitimately sits on a `gd` for longer than anyone will wait, // and the message row is already narrating why. const init_ms = 8_000; // spawn + initialize handshake const req_ms = 4_000; // one request/response round trip const diag_ms = 1_200; // wait for a publishDiagnostics push after didChange const reply_ms = 2_000; // our answers to server-to-client requests const wedged_strikes = 3; // consecutive timeouts before a restart const max_rows = 2000; const max_doc_bytes = 8 << 20; const Err = error{ Dead, Timeout, Protocol, OutOfMemory, NoServer, WriteFailed }; /// Long-lived state outlives every query arena and cannot borrow the caller's /// gpa (a different one shows up in the harness than in the shell), so /// connections own their memory from the page allocator. const sa = std.heap.page_allocator; /// Which units `character` counts in — see the header note on believing the /// server. const Enc = enum { utf8, utf16 }; const Doc = struct { uri: []u8, version: u32, /// content hash: a query whose buffer has not moved since the last one /// sends no didChange at all, which is most of what makes warm queries fast hash: u64, /// a didChange the server has not answered with diagnostics yet stale: bool = true, }; /// The last publishDiagnostics per file, kept as raw params and re-parsed /// against a query's arena. The accumulation IS workspace diagnostics for a /// server with no pull support. const DiagSet = struct { uri: []u8, body: []u8 }; /// The slice of server capabilities this client changes behaviour on. Silent /// kinds (a server with no renameProvider) need no flag — the request errors /// and errors render as no rows. These four either pick between two code /// paths or gate a second round trip. const Caps = struct { enc: Enc = .utf16, /// textDocument/diagnostic (LSP 3.17 pull) — preferred over the push store pull: bool = false, /// workspace/diagnostic pull_workspace: bool = false, call_hier: bool = false, type_hier: bool = false, /// workspace/didChangeWorkspaceFolders is worth sending folders: bool = false, }; const State = enum(u8) { /// never spawned — the row every spec starts on off, /// spawned, initialize in flight starting, ready, /// transport broke or the server was declared wedged; next query respawns dead, /// binary missing or two failed handshakes; stays down for the session disabled, }; const Conn = struct { mu: libc.pthread_mutex_t = .{}, cond: libc.pthread_cond_t = .{}, state: State = .off, /// bumped per spawn. A reader thread that sees a different gen than its /// own is reading a corpse and exits; a waiter that sees one stops waiting. gen: u32 = 0, handshake_fails: u8 = 0, /// monotonic ms before which a dead/failed server is not respawned — /// exponential backoff against forking a doomed child per keystroke retry_after_ms: i64 = 0, /// monotonic ms when the last handshake completed — a crash shortly /// after "ready" counts as a handshake failure for backoff purposes ready_at_ms: i64 = 0, /// linux: the zeroed condvar defaults to REALTIME; re-initialized with a /// monotonic condattr before its first wait (under the mutex) cond_monotonic: bool = false, timeouts: u8 = 0, pid: libc.pid_t = -1, sock: c_int = -1, next_id: u32 = 1, caps: Caps = .{}, /// the workspace root sent in initialize (sa-owned) root: []u8 = &.{}, /// roots added since, via didChangeWorkspaceFolders (sa-owned entries) extra_roots: std.ArrayList([]u8) = .empty, docs: std.ArrayList(Doc) = .empty, diags: std.ArrayList(DiagSet) = .empty, /// the response mailbox: one request outstanding per connection want_id: u32 = 0, resp: ?[]u8 = null, /// active $/progress begins, and the last begin's title for report rows progress: i32 = 0, title: [48]u8 = @splat(0), title_len: u8 = 0, fn lock(c: *Conn) void { _ = libc.pthread_mutex_lock(&c.mu); } fn unlock(c: *Conn) void { _ = libc.pthread_mutex_unlock(&c.mu); } fn alive(c: *const Conn) bool { return c.state == .starting or c.state == .ready; } }; var conns: [specs.len]Conn = @splat(.{}); // ------------------------------------------------------------- status sink /// One registered listener for unsolicited state changes; both native shells /// register at startup and DEREGISTER before tearing their loop down — the /// registry lock is held across the callback, so a null-ing shell cannot race /// a reader thread mid-post. The callback must copy `text` before returning. var sink_mu: std.atomic.Mutex = .unlocked; var sink_ctx: ?*anyopaque = null; var sink_cb: ?*const fn (ctx: ?*anyopaque, text: []const u8) void = null; /// throttle for chatty progress reports, per spec var sink_last_ms: [specs.len]i64 = @splat(0); var sink_last_text: [specs.len][96]u8 = @splat(@splat(0)); var sink_last_len: [specs.len]u8 = @splat(0); pub fn setStatusSink(ctx: ?*anyopaque, cb: ?*const fn (ctx: ?*anyopaque, text: []const u8) void) void { while (!sink_mu.tryLock()) std.atomic.spinLoopHint(); defer sink_mu.unlock(); sink_ctx = ctx; sink_cb = cb; } const Chat = enum { /// a progress report: at most one per 150ms per server, dropped when it /// repeats the previous text — rust-analyzer emits thousands over a big /// index and the message row repaints per post chatty, /// a state change: starting, ready, exited, errors. Always posted. always, }; fn post(si: usize, chat: Chat, comptime fmt: []const u8, args: anytype) void { var buf: [192]u8 = undefined; const text = std.fmt.bufPrint(&buf, fmt, args) catch return; while (!sink_mu.tryLock()) std.atomic.spinLoopHint(); defer sink_mu.unlock(); const cb = sink_cb orelse return; // Repeating the row that is already showing is never news, whatever the // class — and it is what makes the FINAL state deterministic for the // snapshot harness: however many intermediate reports the throttle let // through, an `.always` end state lands exactly once. const cut = @min(text.len, sink_last_text[si].len); if (std.mem.eql(u8, text[0..cut], sink_last_text[si][0..sink_last_len[si]])) return; if (chat == .chatty) { const now = nowMs(); if (now - sink_last_ms[si] < 150) return; sink_last_ms[si] = now; } sink_last_len[si] = @intCast(cut); @memcpy(sink_last_text[si][0..cut], text[0..cut]); cb(sink_ctx, text); } // ------------------------------------------------------------- introspection const LogEntry = struct { used: bool = false, kind: lsp.Kind = .definition, spec: u8 = 0, us: u64 = 0, rows: usize = 0, err: [24]u8 = @splat(0), err_len: u8 = 0, file: [64]u8 = @splat(0), file_len: u8 = 0, }; const log_cap = 24; var log_buf: [log_cap]LogEntry = @splat(.{}); var log_next: usize = 0; var log_total: u64 = 0; var log_mu: std.atomic.Mutex = .unlocked; fn record(si: usize, req: lsp.Req, us: u64, rows: usize, err: []const u8) void { while (!log_mu.tryLock()) std.atomic.spinLoopHint(); defer log_mu.unlock(); const e = &log_buf[log_next]; e.* = .{ .used = true, .kind = req.kind, .spec = @intCast(si), .us = us, .rows = rows }; const base = std.fs.path.basename(req.path); e.file_len = @intCast(@min(base.len, e.file.len)); @memcpy(e.file[0..e.file_len], base[0..e.file_len]); e.err_len = @intCast(@min(err.len, e.err.len)); @memcpy(e.err[0..e.err_len], err[0..e.err_len]); log_next = (log_next + 1) % log_cap; log_total += 1; } fn hideTime() bool { const v = libc.getenv("PARDES_NOTIME") orelse return false; return std.mem.span(v).len != 0; } /// `SPC l w` narration, same shape as the ZLS backend's: threaded through the /// real path, so it cannot disagree with what `gd` actually did. const Trace = struct { on: bool = false, buf: [16 * 1024]u8 = undefined, len: usize = 0, fn note(t: *Trace, comptime fmt: []const u8, args: anytype) void { if (!t.on or t.len == t.buf.len) return; var w: std.Io.Writer = .fixed(t.buf[t.len..]); w.print(fmt ++ "\n", args) catch {}; t.len += w.buffered().len; } }; // ------------------------------------------------------------------- query pub fn query(gpa: std.mem.Allocator, arena: std.mem.Allocator, req: lsp.Req, out: *std.Io.Writer) !void { _ = gpa; if (req.kind == .status) return status(req, out); var tr: Trace = .{ .on = req.kind == .explain }; const si = specFor(req.path) orelse { tr.note("STOP: no language server spec matches {s}", .{ if (req.path.len == 0) "a pane with no file" else req.path, }); return traceOut(&tr, req, out, 0, 0); }; tr.note("file {s} -> {s} (languageId {s})", .{ std.fs.path.basename(req.path), specs[si].name, specs[si].lang }); const scratch_buf = try arena.alloc(u8, max_rows * 512); var scratch: std.Io.Writer = .fixed(scratch_buf); const t0 = nowUs(); var err_name: []const u8 = ""; var failure: ?anyerror = null; answer(arena, si, req, &scratch, &tr) catch |e| { failure = e; err_name = @errorName(e); tr.note("ERROR: {s}", .{err_name}); }; const us = nowUs() -| t0; const rows = std.mem.count(u8, scratch.buffered(), "\n"); record(si, req, us, rows, err_name); if (req.kind == .explain) return traceOut(&tr, req, out, rows, us); if (failure) |err| return err; try out.writeAll(scratch.buffered()); } fn traceOut(tr: *const Trace, req: lsp.Req, out: *std.Io.Writer, rows: usize, us: u64) std.Io.Writer.Error!void { if (req.kind != .explain) return; try out.print("lsp explain — the definition query at byte {d} of {s}\n\n", .{ req.offset, if (req.path.len == 0) "(no file)" else std.fs.path.basename(req.path), }); try out.writeAll(tr.buf[0..tr.len]); if (hideTime()) try out.print("\n{d} row(s)\n", .{rows}) else try out.print("\n{d} row(s) in {d}us\n", .{ rows, us }); } fn answer(arena: std.mem.Allocator, si: usize, req: lsp.Req, out: *std.Io.Writer, tr: *Trace) Err!void { const c = &conns[si]; c.lock(); defer c.unlock(); const g = try ensure(c, si, arena, req, tr); run(c, si, arena, req, out, tr) catch |e| { switch (e) { // transport-level: this server is gone; forget it, and the next // query respawns. Gen-checked so a respawn that happened while we // waited is not the one we kill. error.Dead, error.Protocol => shutdownIf(c, g), error.Timeout => { // Slow is only wedged when the server is IDLE. One with // active `$/progress` work (rust-analyzer running cargo // check over a thousand crates) is demonstrably alive, is // narrating its own excuse on the message row, and killing // it would throw the index away right before it pays off. if (c.progress > 0) { c.timeouts = 0; } else { c.timeouts +|= 1; if (c.timeouts >= wedged_strikes and c.gen == g) { post(si, .always, "{s} not answering — restarting", .{specs[si].name}); shutdownIf(c, g); } } }, error.OutOfMemory, error.NoServer, error.WriteFailed => {}, } return e; }; c.timeouts = 0; } fn run(c: *Conn, si: usize, arena: std.mem.Allocator, req: lsp.Req, out: *std.Io.Writer, tr: *Trace) Err!void { // `explain` narrates the definition query — dispatch below on the // effective kind so the trace follows the code `gd` really runs. const kind: lsp.Kind = if (req.kind == .explain) .definition else req.kind; // Both arg-taking kinds are useless without one, and an empty // workspace/symbol query means "every symbol in the project". if ((kind == .rename or kind == .workspace_symbols) and req.arg.len == 0) return; var uri: std.ArrayList(u8) = .empty; try uriOf(&uri, arena, req.path); try syncDoc(c, si, arena, uri.items, req.source); const pos = posOf(req.source, req.offset, c.caps.enc); var cx: Cx = .{ .arena = arena, .base = req.root, .cur_path = req.path, .cur_src = req.source, .out = out }; tr.note("server {s} (pid {d}), root {s}, {s} columns", .{ @tagName(c.state), c.pid, c.root, @tagName(c.caps.enc), }); const deadline = nowMs() + req_ms; switch (kind) { .definition, .declaration, .type_definition, .implementation, .references => { const method = switch (kind) { .definition => "textDocument/definition", .declaration => "textDocument/declaration", .type_definition => "textDocument/typeDefinition", .implementation => "textDocument/implementation", else => "textDocument/references", }; var b = try atPos(arena, uri.items, pos); if (kind == .references) try app(&b, arena, ",\"context\":{\"includeDeclaration\":true}"); tr.note("-> {s} @ {d}:{d}", .{ method, pos.line + 1, pos.ch + 1 }); const result = (try call(c, arena, method, b.items, deadline)) orelse { tr.note("<- empty (no result, an error response, or null)", .{}); return; }; try locations(&cx, c.caps.enc, result); tr.note("<- {d} row(s)", .{cx.rows}); }, .select_refs => { const b = try atPos(arena, uri.items, pos); const result = (try call(c, arena, "textDocument/documentHighlight", b.items, deadline)) orelse return; for (items(result)) |h| try emitRange(&cx, c.caps.enc, uri.items, get(h, "range"), ""); }, .hover => { const b = try atPos(arena, uri.items, pos); const result = (try call(c, arena, "textDocument/hover", b.items, deadline)) orelse return; const contents = get(result, "contents") orelse return; const text = str(contents) orelse str(get(contents, "value")) orelse blk: { for (items(contents)) |m| if (str(m) orelse str(get(m, "value"))) |t| break :blk t; break :blk ""; }; if (text.len == 0) return; try out.writeAll(std.mem.trim(u8, text, " \t\r\n")); try out.writeByte('\n'); }, .document_symbols => { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "\"textDocument\":{\"uri\":"); try jstr(&b, arena, uri.items); try app(&b, arena, "}"); const result = (try call(c, arena, "textDocument/documentSymbol", b.items, deadline)) orelse return; try walkSymbols(&cx, c.caps.enc, uri.items, result, 0); }, .workspace_symbols => { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "\"query\":"); try jstr(&b, arena, req.arg); const result = (try call(c, arena, "workspace/symbol", b.items, deadline)) orelse return; for (items(result)) |sym| { const loc = get(sym, "location") orelse continue; try emitRange(&cx, c.caps.enc, str(get(loc, "uri")) orelse continue, get(loc, "range"), str(get(sym, "name")) orelse ""); } }, .diagnostics => try diagnostics(c, arena, &cx, uri.items, deadline), .workspace_diagnostics => try workspaceDiagnostics(c, arena, &cx, deadline), .rename => { var b = try atPos(arena, uri.items, pos); try app(&b, arena, ",\"newName\":"); try jstr(&b, arena, req.arg); const result = (try call(c, arena, "textDocument/rename", b.items, deadline)) orelse return; try renameEdits(&cx, c.caps.enc, uri.items, result, req.source); }, .format => { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "\"textDocument\":{\"uri\":"); try jstr(&b, arena, uri.items); try app(&b, arena, "},\"options\":{\"tabSize\":4,\"insertSpaces\":true}"); const result = (try call(c, arena, "textDocument/formatting", b.items, deadline)) orelse return; try formatEdits(&cx, c.caps.enc, result, req.source); }, .code_action => { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "\"textDocument\":{\"uri\":"); try jstr(&b, arena, uri.items); try b.print(arena, "}},\"range\":{{\"start\":{{\"line\":{d},\"character\":{d}}},\"end\":{{\"line\":{d},\"character\":{d}}}}},\"context\":{{\"diagnostics\":[]}}", .{ pos.line, pos.ch, pos.line, pos.ch }); const result = (try call(c, arena, "textDocument/codeAction", b.items, deadline)) orelse return; for (items(result)) |ca| { const title = str(get(ca, "title")) orelse continue; try lsp.row(cx.out, lsp.rel(cx.base, cx.cur_path), pos.line, 0, flat(arena, title)); cx.rows += 1; } }, .completion => { var b = try atPos(arena, uri.items, pos); try app(&b, arena, ",\"context\":{\"triggerKind\":1}"); const result = (try call(c, arena, "textDocument/completion", b.items, deadline)) orelse return; const list = if (get(result, "items")) |it| items(it) else items(result); const here = lsp.rel(cx.base, cx.cur_path); var n: usize = 0; for (list) |item| { if (n >= 100) break; const label = str(get(item, "label")) orelse continue; const detail = str(get(item, "detail")) orelse ""; var text: std.ArrayList(u8) = .empty; try app(&text, arena, label); if (detail.len > 0) { try app(&text, arena, " "); try app(&text, arena, flat(arena, detail)); } // Every row carries the ASKING position: a completion item has // no location of its own (unlike the ZLS backend, which points // at declarations), so the honest place is where it would be // inserted. n/N still step the list; Enter goes nowhere new. try lsp.row(cx.out, here, pos.line, byteCol(req.source, pos.line, pos.ch, c.caps.enc), text.items); n += 1; } }, .incoming_calls => try hierarchy(c, si, arena, &cx, uri.items, pos, deadline, .incoming, tr), .outgoing_calls => try hierarchy(c, si, arena, &cx, uri.items, pos, deadline, .outgoing, tr), .supertypes => try hierarchy(c, si, arena, &cx, uri.items, pos, deadline, .supers, tr), .subtypes => try hierarchy(c, si, arena, &cx, uri.items, pos, deadline, .subs, tr), .status, .explain => unreachable, } } // -------------------------------------------------------- kind sub-handlers /// Goto/references result shapes: bare Location, Location[], LocationLink[]. fn locations(cx: *Cx, enc: Enc, result: std.json.Value) Err!void { if (result == .object) { try emitRange(cx, enc, str(get(result, "uri")) orelse return, get(result, "range"), ""); return; } for (items(result)) |loc| { if (get(loc, "targetUri")) |tu| { const r = get(loc, "targetSelectionRange") orelse get(loc, "targetRange"); try emitRange(cx, enc, str(tu) orelse continue, r, ""); } else { try emitRange(cx, enc, str(get(loc, "uri")) orelse continue, get(loc, "range"), ""); } } } /// DocumentSymbol[] nests (`children`), SymbolInformation[] is flat. fn walkSymbols(cx: *Cx, enc: Enc, uri: []const u8, node: std.json.Value, depth: u8) Err!void { if (depth > 8) return; for (items(node)) |sym| { const name = str(get(sym, "name")) orelse continue; if (get(get(sym, "location"), "range")) |r| { try emitRange(cx, enc, str(get(get(sym, "location"), "uri")) orelse uri, r, name); } else { try emitRange(cx, enc, uri, get(sym, "selectionRange") orelse get(sym, "range"), name); } if (get(sym, "children")) |kids| try walkSymbols(cx, enc, uri, kids, depth + 1); } } /// Pull when the server does (LSP 3.17), the push store otherwise. The store /// path waits briefly for a publish that postdates the didChange we just /// sent, so `]d` right after an edit sees the new truth, not the old one. fn diagnostics(c: *Conn, arena: std.mem.Allocator, cx: *Cx, uri: []const u8, deadline: i64) Err!void { if (c.caps.pull) { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "\"textDocument\":{\"uri\":"); try jstr(&b, arena, uri); try app(&b, arena, "}"); const result = (try call(c, arena, "textDocument/diagnostic", b.items, deadline)) orelse return; for (items(get(result, "items"))) |dg| try emitDiag(cx, c.caps.enc, uri, dg, arena); return; } var stale = true; for (c.docs.items) |d| if (std.mem.eql(u8, d.uri, uri)) { stale = d.stale; break; }; if (stale) waitFresh(c, uri, nowMs() + diag_ms); try renderStore(c, arena, cx, uri); } fn workspaceDiagnostics(c: *Conn, arena: std.mem.Allocator, cx: *Cx, deadline: i64) Err!void { if (c.caps.pull_workspace) { const result = (try call(c, arena, "workspace/diagnostic", "\"previousResultIds\":[]", deadline)) orelse return; for (items(get(result, "items"))) |per| { const uri = str(get(per, "uri")) orelse continue; for (items(get(per, "items"))) |dg| try emitDiag(cx, c.caps.enc, uri, dg, arena); } return; } try renderStore(c, arena, cx, null); } /// One diagnostic row: `severity: message`, at the diagnostic's own range. fn emitDiag(cx: *Cx, enc: Enc, uri: []const u8, dg: std.json.Value, arena: std.mem.Allocator) Err!void { const sev = num(get(dg, "severity")) orelse 1; const label: []const u8 = switch (sev) { 1 => "error", 2 => "warning", 3 => "info", else => "hint", }; const msg = try std.fmt.allocPrint(arena, "{s}: {s}", .{ label, flat(arena, str(get(dg, "message")) orelse "") }); try emitRange(cx, enc, uri, get(dg, "range"), msg); } fn renderStore(c: *Conn, arena: std.mem.Allocator, cx: *Cx, only_uri: ?[]const u8) Err!void { for (c.diags.items) |d| { if (only_uri) |u| if (!std.mem.eql(u8, d.uri, u)) continue; const v = std.json.parseFromSliceLeaky(std.json.Value, arena, d.body, .{}) catch continue; for (items(get(get(v, "params"), "diagnostics"))) |dg| try emitDiag(cx, c.caps.enc, d.uri, dg, arena); } } /// Block (mutex released) until the reader marks `uri` fresh or the deadline /// passes. The reader broadcasts on every publishDiagnostics. fn waitFresh(c: *Conn, uri: []const u8, deadline: i64) void { const g = c.gen; while (c.gen == g and c.alive()) { var fresh = false; for (c.docs.items) |d| if (std.mem.eql(u8, d.uri, uri)) { fresh = !d.stale; break; }; if (fresh) return; if (!timedWait(c, deadline)) return; } } /// A WorkspaceEdit that stays inside the asked-about file becomes `@put` /// records the core applies as one undo transaction; anything wider (a real /// multi-file rename, file creates/renames) becomes a PREVIEW — one location /// row per would-be edit, in the same buffer `gr` fills, because silently /// applying a fraction of a workspace rename would be worse than either. fn renameEdits(cx: *Cx, enc: Enc, self_uri: []const u8, result: std.json.Value, src: []const u8) Err!void { var edits: std.ArrayList(PutEdit) = .empty; var foreign = false; if (get(result, "documentChanges")) |dcs| { for (items(dcs)) |dc| { if (get(dc, "kind") != null) { foreign = true; // create/rename/delete file operations continue; } const u = str(get(get(dc, "textDocument"), "uri")) orelse continue; const in_self = std.mem.eql(u8, u, self_uri); if (!in_self) foreign = true; for (items(get(dc, "edits"))) |ed| { if (in_self) { // one malformed edit poisons the WHOLE mutating response: // applying the valid remainder would be a partial edit set const span = byteSpan(src, get(ed, "range"), enc) orelse return; const text = str(get(ed, "newText")) orelse return; edits.append(cx.arena, .{ .start = span.start, .end = span.end, .text = text }) catch return error.OutOfMemory; } else try emitRange(cx, enc, u, get(ed, "range"), flat(cx.arena, str(get(ed, "newText")) orelse "")); } } } else if (get(result, "changes")) |ch| if (ch == .object) { var it = ch.object.iterator(); while (it.next()) |e| { const in_self = std.mem.eql(u8, e.key_ptr.*, self_uri); if (!in_self) foreign = true; for (items(e.value_ptr.*)) |ed| { if (in_self) { const span = byteSpan(src, get(ed, "range"), enc) orelse return; const text = str(get(ed, "newText")) orelse return; edits.append(cx.arena, .{ .start = span.start, .end = span.end, .text = text }) catch return error.OutOfMemory; } else try emitRange(cx, enc, e.key_ptr.*, get(ed, "range"), flat(cx.arena, str(get(ed, "newText")) orelse "")); } } }; if (foreign) { // the preview needs the self-file rows too — the point is the full map for (edits.items) |ed| { const lc = lsp.lineCol(src, ed.start); try lsp.row(cx.out, lsp.rel(cx.base, cx.cur_path), lc.line, lc.col, flat(cx.arena, ed.text)); cx.rows += 1; } return; } sortEdits(edits.items); for (edits.items) |ed| try lsp.put(cx.out, ed.start, ed.end, ed.text); } /// TextEdit[] from formatting is by definition about the current document: /// straight to sorted `@put` records. fn formatEdits(cx: *Cx, enc: Enc, result: std.json.Value, src: []const u8) Err!void { var edits: std.ArrayList(PutEdit) = .empty; for (items(result)) |ed| { // fail the whole response on the first malformed TextEdit — a subset // of a format is not a format const span = byteSpan(src, get(ed, "range"), enc) orelse return; const text = str(get(ed, "newText")) orelse return; edits.append(cx.arena, .{ .start = span.start, .end = span.end, .text = text }) catch return error.OutOfMemory; } sortEdits(edits.items); for (edits.items) |ed| try lsp.put(cx.out, ed.start, ed.end, ed.text); } /// One would-be buffer mutation, on its way to an `@put` record. const PutEdit = struct { start: usize, end: usize, text: []const u8 }; fn sortEdits(edits: []PutEdit) void { std.mem.sort(PutEdit, edits, {}, struct { fn lt(_: void, a: PutEdit, b: PutEdit) bool { return a.start < b.start; } }.lt); } const Hier = enum { incoming, outgoing, supers, subs }; /// The two-step hierarchy kinds: prepare at the cursor, then follow every /// item the server returned (usually one). Gated on the server capability so /// an old server costs zero round trips. fn hierarchy(c: *Conn, si: usize, arena: std.mem.Allocator, cx: *Cx, uri: []const u8, pos: Pos, deadline: i64, h: Hier, tr: *Trace) Err!void { const call_side = h == .incoming or h == .outgoing; if (call_side and !c.caps.call_hier) { tr.note("STOP: {s} does not advertise callHierarchyProvider", .{specs[si].name}); return; } if (!call_side and !c.caps.type_hier) { tr.note("STOP: {s} does not advertise typeHierarchyProvider", .{specs[si].name}); return; } const prepare: []const u8 = if (call_side) "textDocument/prepareCallHierarchy" else "textDocument/prepareTypeHierarchy"; const follow: []const u8 = switch (h) { .incoming => "callHierarchy/incomingCalls", .outgoing => "callHierarchy/outgoingCalls", .supers => "typeHierarchy/supertypes", .subs => "typeHierarchy/subtypes", }; const b = try atPos(arena, uri, pos); const prepared = (try call(c, arena, prepare, b.items, deadline)) orelse return; for (items(prepared)) |item| { // the item goes back VERBATIM — servers hide resolution state in // `data` and a re-serialized subset would come back unresolvable. // But only if it can round-trip: Stringify panics past its nesting // limit, so a server that nests a bomb in `data` gets no rows for // that item, not a dead editor. if (jsonDepth(item, 0) > 96) continue; var body: std.ArrayList(u8) = .empty; try app(&body, arena, "\"item\":"); const item_json = std.json.Stringify.valueAlloc(arena, item, .{}) catch return error.OutOfMemory; try app(&body, arena, item_json); const result = (try call(c, arena, follow, body.items, deadline)) orelse return; for (items(result)) |entry| switch (h) { .incoming => { // one row per CALL SITE, under the caller's name — that is // what n/N want to walk const from = get(entry, "from") orelse continue; const fu = str(get(from, "uri")) orelse continue; const name = str(get(from, "name")) orelse ""; const ranges = items(get(entry, "fromRanges")); if (ranges.len == 0) { try emitRange(cx, c.caps.enc, fu, get(from, "selectionRange"), name); } else for (ranges) |r| try emitRange(cx, c.caps.enc, fu, r, name); }, .outgoing => { const to = get(entry, "to") orelse continue; try emitRange(cx, c.caps.enc, str(get(to, "uri")) orelse continue, get(to, "selectionRange") orelse get(to, "range"), hierText(cx.arena, to)); }, .supers, .subs => try emitRange(cx, c.caps.enc, str(get(entry, "uri")) orelse continue, get(entry, "selectionRange") orelse get(entry, "range"), hierText(cx.arena, entry)), }; } } /// Depth of a parsed json value, saturating just past `at`'s caller's cap — /// the guard that keeps server-controlled nesting away from Stringify's /// fixed-depth assertion. fn jsonDepth(v: std.json.Value, at: u32) u32 { if (at > 96) return at; return switch (v) { .object => |o| blk: { var deepest = at; var it = o.iterator(); while (it.next()) |e| deepest = @max(deepest, jsonDepth(e.value_ptr.*, at + 1)); break :blk deepest; }, .array => |a| blk: { var deepest = at; for (a.items) |e| deepest = @max(deepest, jsonDepth(e, at + 1)); break :blk deepest; }, else => at, }; } fn hierText(arena: std.mem.Allocator, item: std.json.Value) []const u8 { const name = str(get(item, "name")) orelse ""; const detail = str(get(item, "detail")) orelse return name; return std.fmt.allocPrint(arena, "{s} {s}", .{ name, flat(arena, detail) }) catch name; } // ------------------------------------------------------------ row rendering /// Files read while answering ONE query, so a references list over a handful /// of files reads each once, not once per row. The current buffer never needs /// reading: `req.source` IS its text, unsaved edits included. const Cx = struct { arena: std.mem.Allocator, base: []const u8, cur_path: []const u8, cur_src: []const u8, out: *std.Io.Writer, rows: usize = 0, files: std.ArrayList(struct { path: []const u8, text: []const u8 }) = .empty, fn text(cx: *Cx, path: []const u8) ?[]const u8 { if (std.mem.eql(u8, path, cx.cur_path)) return cx.cur_src; for (cx.files.items) |f| if (std.mem.eql(u8, f.path, path)) return f.text; if (cx.files.items.len >= 32) return null; var pathbuf: [4096]u8 = undefined; const path_z = std.fmt.bufPrintSentinel(&pathbuf, "{s}", .{path}, 0) catch return null; const fd = libc.open(path_z, .{ .ACCMODE = .RDONLY }); if (fd < 0) return null; defer _ = libc.close(fd); var buf: std.ArrayList(u8) = .empty; var chunk: [64 * 1024]u8 = undefined; while (buf.items.len < max_doc_bytes) { const n = libc.read(fd, &chunk, chunk.len); if (n < 0) { if (libc.errno(n) == .INTR) continue; return null; } if (n == 0) break; buf.appendSlice(cx.arena, chunk[0..@intCast(n)]) catch return null; } cx.files.append(cx.arena, .{ .path = path, .text = buf.items }) catch return null; return buf.items; } }; /// One row from an LSP (uri, Range). The uri becomes a real path (percent- /// decoded, `rel`'d against the asking window), utf-16 columns become byte /// columns, and a single-line range becomes the `path:LINE:COL-ENDCOL` form a /// look SELECTS. `note` overrides the source line as the row's text. fn emitRange(cx: *Cx, enc: Enc, uri: []const u8, range: ?std.json.Value, note: []const u8) Err!void { if (cx.rows >= max_rows) return; const path = pathOf(cx.arena, uri) orelse return; if (path.len == 0 or path[0] != '/') return; // rows promise absolute-or-rel-from-base const r = rangeOf(range) orelse return; const src = cx.text(path); var bol: usize = 0; var eol: usize = 0; if (src) |t| { var n: u32 = 0; while (n < r.sl) : (n += 1) { bol = (std.mem.indexOfScalarPos(u8, t, bol, '\n') orelse { bol = t.len; break; }) + 1; } eol = std.mem.indexOfScalarPos(u8, t, @min(bol, t.len), '\n') orelse t.len; } const lntext: []const u8 = if (src) |t| t[@min(bol, t.len)..@min(eol, t.len)] else ""; const col = colBytes(lntext, r.sc, enc); const shown = lsp.rel(cx.base, path); const rowtext = if (note.len > 0) note else lntext; if (r.el == r.sl and r.ec > r.sc) { // spanRow wants the protocol's EXCLUSIVE end as a 1-based inclusive // byte column; converting the exclusive utf-16 end unit yields the // exclusive byte column, which is the same number. const end_col = colBytes(lntext, r.ec, enc); try lsp.spanRow(cx.out, shown, r.sl, col, r.el, end_col, rowtext); } else { try lsp.row(cx.out, shown, r.sl, col, rowtext); } cx.rows += 1; } /// utf-16 code units -> byte column within one line; identity for utf-8. fn colBytes(lntext: []const u8, ch: u32, enc: Enc) usize { if (enc == .utf8 or lntext.len == 0) return ch; var units: u32 = 0; var i: usize = 0; while (i < lntext.len and units < ch) { const l = std.unicode.utf8ByteSequenceLength(lntext[i]) catch 1; units += if (l == 4) 2 else 1; i += l; } return i; } /// The current-buffer byte column for a (line, ch) position — used only for /// rows that point at the asking position itself. fn byteCol(src: []const u8, line: u32, ch: u32, enc: Enc) usize { var bol: usize = 0; var n: u32 = 0; while (n < line) : (n += 1) bol = (std.mem.indexOfScalarPos(u8, src, bol, '\n') orelse return ch) + 1; const eol = std.mem.indexOfScalarPos(u8, src, bol, '\n') orelse src.len; return colBytes(src[bol..eol], ch, enc); } /// Whole-document Range -> half-open byte span for a MUTATING record. /// STRICT: null unless both endpoints denote positions that exist in the /// document — for edits, a clamped range is a wrong edit, and wrong edits /// fail closed (the render paths keep their forgiving conversions; a row a /// column off is an inconvenience, a splice a column off is corruption). fn byteSpan(src: []const u8, range: ?std.json.Value, enc: Enc) ?struct { start: usize, end: usize } { const r = rangeOf(range) orelse return null; const start = strictOffset(src, r.sl, r.sc, enc) orelse return null; const end = strictOffset(src, r.el, r.ec, enc) orelse return null; if (end < start) return null; return .{ .start = start, .end = end }; } /// (line, character) -> byte offset, or null when the line does not exist or /// the character runs past its end. End-of-line (character == line length) is /// a real position — that is where an insert-at-EOL lands. fn strictOffset(src: []const u8, line: u32, ch: u32, enc: Enc) ?usize { var bol: usize = 0; var n: u32 = 0; while (n < line) : (n += 1) bol = (std.mem.indexOfScalarPos(u8, src, bol, '\n') orelse return null) + 1; const eol = std.mem.indexOfScalarPos(u8, src, bol, '\n') orelse src.len; if (enc == .utf8) { if (ch > eol - bol) return null; return bol + ch; } var units: u32 = 0; var i: usize = bol; while (i < eol and units < ch) { const l = std.unicode.utf8ByteSequenceLength(src[i]) catch 1; units += if (l == 4) 2 else 1; i += l; } if (units < ch) return null; return i; } /// The RENDER-path sibling of strictOffset: clamps instead of failing, /// because a row is presentation, not mutation. fn offsetAt(src: []const u8, line: u32, ch: u32, enc: Enc) usize { var bol: usize = 0; var n: u32 = 0; while (n < line) : (n += 1) bol = (std.mem.indexOfScalarPos(u8, src, bol, '\n') orelse return src.len) + 1; const eol = std.mem.indexOfScalarPos(u8, src, bol, '\n') orelse src.len; return bol + @min(colBytes(src[bol..eol], ch, enc), eol - bol); } /// Prose squashed onto the one line a row is: newlines and tabs to spaces, /// cut to a width a pane can show. fn flat(arena: std.mem.Allocator, s: []const u8) []const u8 { var buf: std.ArrayList(u8) = .empty; var sp = false; for (s) |ch| { if (ch == '\n' or ch == '\r' or ch == '\t' or ch == 0) { if (!sp and buf.items.len > 0) buf.append(arena, ' ') catch break; sp = true; } else { buf.append(arena, ch) catch break; sp = false; } if (buf.items.len >= 200) break; } return std.mem.trim(u8, buf.items, " "); } fn setCloexec(fd: c_int) void { const FD_CLOEXEC: c_int = 1; _ = libc.fcntl(fd, libc.F.SETFD, FD_CLOEXEC); } // Darwin needs fcntl after socketpair; another concurrent fork can inherit the pair in that window. fn transportPair(sv: *[2]libc.fd_t) bool { const sock_type = if (comptime builtin.os.tag.isDarwin()) libc.SOCK.STREAM else libc.SOCK.STREAM | libc.SOCK.CLOEXEC; if (libc.socketpair(libc.AF.UNIX, sock_type, 0, sv) != 0) return false; if (comptime builtin.os.tag.isDarwin()) { setCloexec(sv[0]); // The child dup2s this onto 0 and 1, and dup2 CLEARS close-on-exec on // the copy, so the server still gets the socket; this marks only the // number itself, which the child's 3..1024 sweep closes anyway. setCloexec(sv[1]); const one: c_int = 1; _ = libc.setsockopt(sv[0], libc.SOL.SOCKET, so_nosigpipe, @ptrCast(&one), @sizeOf(c_int)); } return true; } // ----------------------------------------------------------- the connection /// Spawn-or-return, and tell an existing server about a new project root. /// Returns the generation the caller's whole query is bound to. fn ensure(c: *Conn, si: usize, arena: std.mem.Allocator, req: lsp.Req, tr: *Trace) Err!u32 { if (c.state == .disabled) { tr.note("STOP: {s} is disabled for this session (no binary; set {s})", .{ specs[si].name, specs[si].env }); return error.NoServer; } if (c.alive()) { const root = rootOf(req.root, specs[si].markers); if (!std.mem.eql(u8, root, c.root)) { var known = false; for (c.extra_roots.items) |r| if (std.mem.eql(u8, r, root)) { known = true; break; }; if (!known and c.caps.folders) { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "{\"jsonrpc\":\"2.0\",\"method\":\"workspace/didChangeWorkspaceFolders\",\"params\":{\"event\":{\"added\":[{\"name\":\"root\",\"uri\":\""); try uriOf(&b, arena, root); try app(&b, arena, "\"}],\"removed\":[]}}}"); try frame(c.sock, b.items, nowMs() + reply_ms); const owned = sa.dupe(u8, root) catch return error.OutOfMemory; c.extra_roots.append(sa, owned) catch sa.free(owned); } } return c.gen; } // A recent failed spawn/handshake holds the fork back; the message row // already said when the next try is due. if (nowMs() < c.retry_after_ms) { tr.note("STOP: {s} failed recently; retry due in {d}ms", .{ specs[si].name, c.retry_after_ms - nowMs() }); return error.NoServer; } // Look before forking; a missing binary must cost one probe and one // message, not a doomed fork per keystroke. var exe_buf: [4096:0]u8 = undefined; const exe = binOf(&specs[si], &exe_buf) orelse { c.state = .disabled; post(si, .always, "{s} not found — install it or set {s}", .{ specs[si].name, specs[si].env }); tr.note("STOP: no {s} binary on PATH (override: {s})", .{ specs[si].bin, specs[si].env }); return error.NoServer; }; const root = rootOf(req.root, specs[si].markers); var root_buf: [4096:0]u8 = undefined; const root_z = std.fmt.bufPrintSentinel(&root_buf, "{s}", .{root}, 0) catch return error.NoServer; // Owned BEFORE anything spawns: the reader thread is the socket's only // closer once it exists, so nothing fallible may sit between fork and // reader-spawn — an abort there would strand the fd. const owned_root = sa.dupe(u8, root) catch return error.OutOfMemory; var sv: [2]libc.fd_t = undefined; if (!transportPair(&sv)) { sa.free(owned_root); tr.note("STOP: socketpair for {s} failed", .{specs[si].name}); return error.NoServer; } const pid = libc.fork(); if (pid < 0) { _ = libc.close(sv[0]); _ = libc.close(sv[1]); sa.free(owned_root); tr.note("STOP: fork for {s} failed", .{specs[si].name}); return error.NoServer; } if (pid == 0) { // Between fork and exec, in a process with threads: only async- // signal-safe calls, no allocation, no locks. (Same rule as tty.zig's // forkShell.) _ = setsid(); // Ctrl-C in pardes's terminal is not the server's business _ = libc.dup2(sv[1], 0); _ = libc.dup2(sv[1], 1); const devnull = libc.open("/dev/null", .{ .ACCMODE = .WRONLY }); if (devnull >= 0) _ = libc.dup2(devnull, 2); var fd: c_int = 3; while (fd < 1024) : (fd += 1) _ = libc.close(fd); _ = chdir(root_z.ptr); var argv: [8:null]?[*:0]const u8 = @splat(null); argv[0] = exe.ptr; // args are comptime literals; the buffers live until execv var argbufs: [6][64:0]u8 = undefined; for (specs[si].args, 0..) |a, i| { if (i >= argbufs.len) break; const z = std.fmt.bufPrintSentinel(&argbufs[i], "{s}", .{a}, 0) catch break; argv[1 + i] = z.ptr; } _ = execv(exe.ptr, &argv); _exit(127); } _ = libc.close(sv[1]); // the slate the old generation may have left (a self-died server skips // shutdownIf) must not leak into the new one: a stale doc entry would // suppress the didOpen the new server never got c.gen +%= 1; const g = c.gen; c.pid = pid; c.sock = sv[0]; c.state = .starting; c.caps = .{}; c.next_id = 1; c.want_id = 0; if (c.resp) |r| sa.free(r); c.resp = null; c.progress = 0; c.timeouts = 0; forgetDocs(c); if (c.root.len > 0) sa.free(c.root); c.root = owned_root; const th = std.Thread.spawn(.{}, reader, .{ si, g, sv[0] }) catch { // no reader means nobody would ever close the fd our way; do it here, // before anything else can see the conn _ = libc.close(sv[0]); _ = libc.kill(pid, .KILL); _ = libc.waitpid(pid, null, 0); c.state = .dead; c.sock = -1; return error.NoServer; }; th.detach(); post(si, .always, "{s} starting — {s}", .{ specs[si].name, root }); tr.note("spawned {s} (pid {d}) at {s}", .{ specs[si].name, pid, root }); handshake(c, arena, root) catch |e| { shutdownIf(c, g); // NOT a session disable. A handshake that misses the deadline is // routinely environmental — a rustup shim deciding to download the // project's pinned toolchain before launching the real server was // the case that taught this — and it heals by itself. What must not // happen is a fork per keystroke while it heals, so failures back // off on backoffMs's schedule, reset by the next success. c.handshake_fails +|= 1; const wait = backoffMs(c.handshake_fails); c.retry_after_ms = nowMs() + wait; post(si, .always, "{s} not answering — retrying in {d}s", .{ specs[si].name, @divTrunc(wait, 1000), }); return switch (e) { error.OutOfMemory => error.OutOfMemory, else => error.NoServer, }; }; c.handshake_fails = 0; c.retry_after_ms = 0; c.ready_at_ms = nowMs(); c.state = .ready; post(si, .always, "{s} ready — {s}", .{ specs[si].name, root }); return g; } fn handshake(c: *Conn, arena: std.mem.Allocator, root: []const u8) Err!void { var b: std.ArrayList(u8) = .empty; try b.print(arena, "\"processId\":{d},\"clientInfo\":{{\"name\":\"pardes\"}},\"rootUri\":\"", .{@as(u32, @bitCast(libc.getpid()))}); try uriOf(&b, arena, root); try app(&b, arena, "\",\"workspaceFolders\":[{\"name\":\"root\",\"uri\":\""); try uriOf(&b, arena, root); try app(&b, arena, "\""); try app(&b, arena, \\}],"capabilities":{"general":{"positionEncodings":["utf-8","utf-16"]}, \\"window":{"workDoneProgress":true}, \\"workspace":{"configuration":true,"workspaceFolders":true,"symbol":{}, \\"diagnostics":{"refreshSupport":false}}, \\"textDocument":{"synchronization":{"dynamicRegistration":false,"willSave":false,"didSave":false}, \\"publishDiagnostics":{"relatedInformation":false}, \\"diagnostic":{"dynamicRegistration":false,"relatedDocumentSupport":false}, \\"hover":{"contentFormat":["plaintext","markdown"]}, \\"definition":{"linkSupport":true},"declaration":{"linkSupport":true}, \\"typeDefinition":{"linkSupport":true},"implementation":{"linkSupport":true}, \\"references":{},"documentHighlight":{}, \\"documentSymbol":{"hierarchicalDocumentSymbolSupport":true}, \\"formatting":{},"rename":{"prepareSupport":false}, \\"completion":{"completionItem":{"snippetSupport":false,"documentationFormat":["plaintext"]}}, \\"callHierarchy":{},"typeHierarchy":{}, \\"codeAction":{"codeActionLiteralSupport":{"codeActionKind":{"valueSet":[]}}}}, \\"experimental":{"serverStatusNotification":true}} ); // strip the literal's newlines: legal JSON either way, but the frame // length must match what is sent var body: std.ArrayList(u8) = .empty; for (b.items) |ch| if (ch != '\n') try body.append(arena, ch); const deadline = nowMs() + init_ms; const reply = (try call(c, arena, "initialize", body.items, deadline)) orelse return error.Protocol; const caps = get(reply, "capabilities"); c.caps.enc = if (std.mem.eql(u8, str(get(caps, "positionEncoding")) orelse "utf-16", "utf-8")) .utf8 else .utf16; if (get(caps, "diagnosticProvider")) |dp| { c.caps.pull = provider(dp); c.caps.pull_workspace = if (get(dp, "workspaceDiagnostics")) |w| w == .bool and w.bool else false; } c.caps.call_hier = provider(get(caps, "callHierarchyProvider")); c.caps.type_hier = provider(get(caps, "typeHierarchyProvider")); c.caps.folders = if (get(get(get(caps, "workspace"), "workspaceFolders"), "supported")) |s| s == .bool and s.bool else false; try frame(c.sock, "{\"jsonrpc\":\"2.0\",\"method\":\"initialized\",\"params\":{}}", deadline); } /// A server capability that may be `true`, an options object, or absent. fn provider(v: ?std.json.Value) bool { const o = v orelse return false; return switch (o) { .bool => |b| b, .object => true, else => false, }; } /// didOpen the first time a file is seen, didChange when its bytes moved, /// nothing when they did not — the common case between two presses of gd. fn syncDoc(c: *Conn, si: usize, arena: std.mem.Allocator, uri: []const u8, src: []const u8) Err!void { const hash = std.hash.Wyhash.hash(0, src); var doc: ?*Doc = null; for (c.docs.items) |*d| if (std.mem.eql(u8, d.uri, uri)) { doc = d; break; }; if (doc) |d| if (d.hash == hash) return; var b: std.ArrayList(u8) = .empty; if (doc) |d| { d.version += 1; d.hash = hash; d.stale = true; try b.print(arena, "{{\"jsonrpc\":\"2.0\",\"method\":\"textDocument/didChange\",\"params\":{{\"textDocument\":{{\"uri\":", .{}); try jstr(&b, arena, uri); try b.print(arena, ",\"version\":{d}}},\"contentChanges\":[{{\"text\":", .{d.version}); try jstr(&b, arena, src); try app(&b, arena, "}]}}"); } else { if (c.docs.items.len >= 256) return; // a session does not open this many const owned = sa.dupe(u8, uri) catch return error.OutOfMemory; c.docs.append(sa, .{ .uri = owned, .version = 1, .hash = hash }) catch { sa.free(owned); return error.OutOfMemory; }; try app(&b, arena, "{\"jsonrpc\":\"2.0\",\"method\":\"textDocument/didOpen\",\"params\":{\"textDocument\":{\"uri\":"); try jstr(&b, arena, uri); try b.print(arena, ",\"languageId\":\"{s}\",\"version\":1,\"text\":", .{specs[si].lang}); try jstr(&b, arena, src); try app(&b, arena, "}}}"); } try frame(c.sock, b.items, nowMs() + req_ms); } /// Send a request and wait on the mailbox. Returns the parsed `result`, or /// null for an error response / null result — both are legal "no answer". /// Runs with the conn mutex held; the mutex is released inside the wait. fn call(c: *Conn, arena: std.mem.Allocator, method: []const u8, params: []const u8, deadline: i64) Err!?std.json.Value { const g = c.gen; const id = c.next_id; c.next_id +%= 1; var b: std.ArrayList(u8) = .empty; try b.print(arena, "{{\"jsonrpc\":\"2.0\",\"id\":{d},\"method\":\"{s}\",\"params\":{{", .{ id, method }); try app(&b, arena, params); try app(&b, arena, "}}"); if (c.resp) |r| sa.free(r); c.resp = null; c.want_id = id; defer c.want_id = 0; try frame(c.sock, b.items, deadline); while (c.resp == null and c.gen == g and c.alive()) { if (!timedWait(c, deadline)) break; } if (c.gen != g or !c.alive()) return error.Dead; const raw = c.resp orelse { // give the server leave to abandon the work we stopped waiting for var cb: std.ArrayList(u8) = .empty; cb.print(arena, "{{\"jsonrpc\":\"2.0\",\"method\":\"$/cancelRequest\",\"params\":{{\"id\":{d}}}}}", .{id}) catch return error.Timeout; frame(c.sock, cb.items, nowMs() + reply_ms) catch {}; return error.Timeout; }; c.resp = null; defer sa.free(raw); const v = std.json.parseFromSliceLeaky(std.json.Value, arena, raw, .{}) catch return error.Protocol; const result = get(v, "result") orelse return null; if (result == .null) return null; return result; } /// One bounded cond wait. False when the deadline has passed. /// /// The deadline arithmetic everywhere here is CLOCK_MONOTONIC, so the wait /// must be too: a REALTIME abstime plus an NTP step backwards would park a /// worker far past `req_ms`, violating the seam's bounded-wait promise. On /// linux the condvar is lazily initialized with a monotonic condattr; darwin /// has no setclock but has a RELATIVE wait, which never consults the wall /// clock at all. fn timedWait(c: *Conn, deadline: i64) bool { const left = deadline - nowMs(); if (left <= 0) return false; const wait_ms = @min(left, 500); // tick so gen/state changes are noticed if (comptime builtin.os.tag.isDarwin()) { const rel: libc.timespec = .{ .sec = @divTrunc(wait_ms, 1000), .nsec = @rem(wait_ms, 1000) * 1_000_000, }; _ = pthread_cond_timedwait_relative_np(&c.cond, &c.mu, &rel); return true; } if (!c.cond_monotonic) { // lazily re-initialize the zeroed (REALTIME) condvar with a // monotonic clock; under the mutex, and only before the first wait, // so no thread can be parked on the old one var attr: pthread_condattr_t = undefined; if (pthread_condattr_init(&attr) == 0) { _ = pthread_condattr_setclock(&attr, CLOCK_MONOTONIC); _ = pthread_cond_init(&c.cond, &attr); _ = pthread_condattr_destroy(&attr); } c.cond_monotonic = true; } var now: libc.timespec = undefined; _ = libc.clock_gettime(.MONOTONIC, &now); var abs: libc.timespec = .{ .sec = now.sec + @divTrunc(wait_ms, 1000), .nsec = now.nsec + @rem(wait_ms, 1000) * 1_000_000, }; if (abs.nsec >= 1_000_000_000) { abs.sec += 1; abs.nsec -= 1_000_000_000; } _ = libc.pthread_cond_timedwait(&c.cond, &c.mu, &abs); return true; } /// Forget a server, but only the generation the caller actually used — a /// respawn that happened during the caller's wait must not be killed for its /// predecessor's crimes. Mutex held. The reader closes the fd; this only /// shuts the transport down and reaps the child. fn shutdownIf(c: *Conn, g: u32) void { if (c.gen != g or c.sock < 0) return; _ = libc.shutdown(c.sock, 2); // SHUT_RDWR: the reader sees EOF and closes reap(c.pid); c.sock = -1; c.pid = -1; if (c.state != .disabled) c.state = .dead; forgetDocs(c); if (c.resp) |r| sa.free(r); c.resp = null; _ = libc.pthread_cond_broadcast(&c.cond); } /// TERM, a bounded grace, then KILL. Never blocks unboundedly: a reaped or /// foreign pid answers -1 immediately. fn reap(pid: libc.pid_t) void { if (pid <= 0) return; _ = libc.kill(pid, .TERM); var tries: u8 = 0; while (tries < 20) : (tries += 1) { const w = libc.waitpid(pid, null, libc.W.NOHANG); if (w == pid or w < 0) return; _ = usleep(10_000); } _ = libc.kill(pid, .KILL); _ = libc.waitpid(pid, null, 0); } /// Free every per-document accumulation. Mutex held. fn forgetDocs(c: *Conn) void { for (c.extra_roots.items) |r| sa.free(r); c.extra_roots.clearRetainingCapacity(); for (c.docs.items) |d| sa.free(d.uri); c.docs.clearRetainingCapacity(); for (c.diags.items) |d| { sa.free(d.uri); sa.free(d.body); } c.diags.clearRetainingCapacity(); } // ------------------------------------------------------------ reader thread /// Owns the read side of one socket for one generation, and is the only /// closer of that fd. Routes responses to the mailbox, answers server /// requests, feeds the diagnostics store, narrates progress. fn reader(si: usize, g: u32, sock: c_int) void { const c = &conns[si]; var buf = sa.alloc(u8, 256 * 1024) catch { readerExit(c, si, g, sock); return; }; defer sa.free(buf); var len: usize = 0; outer: while (true) { // drain every complete frame already buffered drain: while (true) switch (frameNext(buf[0..len])) { .frame => |m| { dispatch(c, si, g, sock, buf[m.body_start..m.body_end]); std.mem.copyForwards(u8, buf[0 .. len - m.consumed], buf[m.consumed..len]); len -= m.consumed; }, .incomplete => break :drain, // a complete header that is not a frame we can speak: nothing // after it can ever re-synchronize, so the connection is over — // this must NOT read as "incomplete", which would buffer the // poison forever while every request quietly times out .poison => break :outer, }; // grow when a frame is bigger than the space that is left if (len == buf.len) { if (buf.len >= 64 << 20) break :outer; const bigger = sa.realloc(buf, buf.len * 2) catch break :outer; buf = bigger; } var pfd = [_]libc.pollfd{.{ .fd = sock, .events = POLLIN, .revents = 0 }}; const pr = libc.poll(&pfd, 1, 1000); if (pr < 0) { if (libc.errno(pr) == .INTR) continue; break :outer; } if (pr == 0) { // tick: a respawn while the server was silent leaves this thread // reading a corpse; notice and go c.lock(); const stale = c.gen != g; c.unlock(); if (stale) break :outer; continue; } const got = libc.recv(sock, buf.ptr + len, buf.len - len, 0); if (got < 0) { const e = libc.errno(got); if (e == .INTR or e == .AGAIN) continue; break :outer; } if (got == 0) break :outer; // the server exited or shutdownIf spoke len += @intCast(got); } readerExit(c, si, g, sock); } fn readerExit(c: *Conn, si: usize, g: u32, sock: c_int) void { c.lock(); if (c.gen == g and c.alive()) { // the server died on its own — shutdownIf never saw it, so the corpse // is ours to reap and the conn's transport fields are ours to clear c.state = .dead; reap(c.pid); c.pid = -1; c.sock = -1; // A server that keeps crashing right after its handshake would // otherwise respawn on every keystroke — the handshake SUCCEEDING // resets the backoff, so the crash has to count as the failure it // is. A long-lived server that dies gets one free respawn. if (nowMs() - c.ready_at_ms < 30_000) { c.handshake_fails +|= 1; const wait = backoffMs(c.handshake_fails); c.retry_after_ms = nowMs() + wait; post(si, .always, "{s} crashed — retrying in {d}s", .{ specs[si].name, @divTrunc(wait, 1000) }); } else { post(si, .always, "{s} exited — restarts on the next query", .{specs[si].name}); } _ = libc.pthread_cond_broadcast(&c.cond); } c.unlock(); _ = libc.close(sock); } /// 10s doubling to 2min: the fork-per-keystroke guard's schedule. fn backoffMs(fails: u8) i64 { const shift: u5 = @min(fails -| 1, 4); return @as(i64, 10_000) << shift; } const Framed = struct { body_start: usize, body_end: usize, consumed: usize }; const FrameStep = union(enum) { /// keep reading; the head of the buffer may still become a frame incomplete, /// a complete, valid frame frame: Framed, /// a complete header that cannot be a frame (no/invalid/oversized /// Content-Length): the stream can never re-synchronize poison, }; /// One Content-Length-framed message at the head of `data`. Pure — this is /// the whole wire format, and it is testable without a socket. fn frameNext(data: []const u8) FrameStep { const sep = std.mem.indexOf(u8, data, "\r\n\r\n") orelse return .incomplete; var clen: ?usize = null; var it = std.mem.splitSequence(u8, data[0..sep], "\r\n"); while (it.next()) |line| { const colon = std.mem.indexOfScalar(u8, line, ':') orelse continue; if (!std.ascii.eqlIgnoreCase(std.mem.trim(u8, line[0..colon], " "), "content-length")) continue; clen = std.fmt.parseInt(usize, std.mem.trim(u8, line[colon + 1 ..], " "), 10) catch null; } const n = clen orelse return .poison; const start = sep + 4; if (n > 64 << 20) return .poison; if (data.len < start + n) return .incomplete; return .{ .frame = .{ .body_start = start, .body_end = start + n, .consumed = start + n } }; } /// One message from the server. Parses into a throwaway arena; touches conn /// state only under the mutex; replies never hold it (bounded writes must not /// stall the mailbox). fn dispatch(c: *Conn, si: usize, g: u32, sock: c_int, raw: []const u8) void { var arena_state: std.heap.ArenaAllocator = .init(sa); defer arena_state.deinit(); const arena = arena_state.allocator(); const v = std.json.parseFromSliceLeaky(std.json.Value, arena, raw, .{}) catch return; if (v != .object) return; if (str(get(v, "method"))) |method| { if (get(v, "id")) |id| return serverRequest(arena, sock, method, id, get(v, "params")); return notification(c, si, g, raw, method, get(v, "params")); } // a response: deliver to the one waiting query, if it is still waiting const id = num(get(v, "id")) orelse return; c.lock(); defer c.unlock(); if (c.gen == g and c.want_id != 0 and id == c.want_id and c.resp == null) { c.resp = sa.dupe(u8, raw) catch null; _ = libc.pthread_cond_broadcast(&c.cond); } } /// The server asked US something. Everything optional was declined in the /// handshake, so a null-ish answer is always legal — but it must ARRIVE, or /// the server blocks forever on its own question. fn serverRequest(arena: std.mem.Allocator, sock: c_int, method: []const u8, id: std.json.Value, params: ?std.json.Value) void { var b: std.ArrayList(u8) = .empty; app(&b, arena, "{\"jsonrpc\":\"2.0\",\"id\":") catch return; switch (id) { .integer => |n| b.print(arena, "{d}", .{n}) catch return, .string => |t| jstr(&b, arena, t) catch return, else => app(&b, arena, "null") catch return, } if (std.mem.eql(u8, method, "workspace/configuration")) { // one null per asked item — a bare null here is a protocol violation // some servers punish with a parse loop const n = items(get(params, "items")).len; app(&b, arena, ",\"result\":[") catch return; for (0..n) |i| app(&b, arena, if (i == 0) "null" else ",null") catch return; app(&b, arena, "]}") catch return; } else if (std.mem.eql(u8, method, "workspace/applyEdit")) { app(&b, arena, ",\"result\":{\"applied\":false}}") catch return; } else { app(&b, arena, ",\"result\":null}") catch return; } frame(sock, b.items, nowMs() + reply_ms) catch {}; } fn notification(c: *Conn, si: usize, g: u32, raw: []const u8, method: []const u8, params: ?std.json.Value) void { if (std.mem.eql(u8, method, "textDocument/publishDiagnostics")) { const uri = str(get(params, "uri")) orelse return; // keep the raw MESSAGE (params included); re-parsing it per query // beats owning a parsed tree past this arena's life const body = sa.dupe(u8, raw) catch return; c.lock(); defer c.unlock(); if (c.gen != g) { sa.free(body); return; } for (c.docs.items) |*d| if (std.mem.eql(u8, d.uri, uri)) { d.stale = false; break; }; var slot: ?*DiagSet = null; for (c.diags.items) |*d| if (std.mem.eql(u8, d.uri, uri)) { slot = d; break; }; if (slot) |d| { sa.free(d.body); d.body = body; } else if (c.diags.items.len < 256) { const owned = sa.dupe(u8, uri) catch { sa.free(body); return; }; c.diags.append(sa, .{ .uri = owned, .body = body }) catch { sa.free(owned); sa.free(body); return; }; } else sa.free(body); _ = libc.pthread_cond_broadcast(&c.cond); // waitFresh watches this return; } if (std.mem.eql(u8, method, "$/progress")) { const value = get(params, "value") orelse return; const pkind = str(get(value, "kind")) orelse return; const msg = str(get(value, "message")) orelse ""; c.lock(); if (c.gen != g) { c.unlock(); return; } if (std.mem.eql(u8, pkind, "begin")) { c.progress += 1; const title = str(get(value, "title")) orelse ""; c.title_len = @intCast(@min(title.len, c.title.len)); @memcpy(c.title[0..c.title_len], title[0..c.title_len]); var tbuf: [48]u8 = undefined; const t = tbuf[0..c.title_len]; @memcpy(t, c.title[0..c.title_len]); c.unlock(); post(si, .chatty, "{s}: {s}\u{2026} {s}", .{ specs[si].name, t, msg }); return; } if (std.mem.eql(u8, pkind, "report")) { var tbuf: [48]u8 = undefined; const t = tbuf[0..c.title_len]; @memcpy(t, c.title[0..c.title_len]); const pct = num(get(value, "percentage")); c.unlock(); if (pct) |p| post(si, .chatty, "{s}: {s} {d}% {s}", .{ specs[si].name, t, p, msg }) else post(si, .chatty, "{s}: {s} {s}", .{ specs[si].name, t, msg }); return; } // end c.progress -= 1; const idle = c.progress <= 0; if (idle) c.progress = 0; c.unlock(); // the setpoint every burst of progress ends on; `.always` + the // dedupe above means it lands exactly once however the reports raced if (idle) post(si, .always, "{s}: ready", .{specs[si].name}); return; } if (std.mem.eql(u8, method, "window/showMessage")) { const t = num(get(params, "type")) orelse 3; if (t > 2) return; // info/log are the server's diary, not the user's post(si, .always, "{s}: {s}", .{ specs[si].name, str(get(params, "message")) orelse "" }); return; } // rust-analyzer's precise quiescence signal, opted into via // experimental.serverStatusNotification if (std.mem.eql(u8, method, "experimental/serverStatus")) { const health = str(get(params, "health")) orelse "ok"; const quiescent = if (get(params, "quiescent")) |q| q == .bool and q.bool else false; if (!std.mem.eql(u8, health, "ok")) post(si, .always, "{s}: {s} — {s}", .{ specs[si].name, health, str(get(params, "message")) orelse "" }) else if (quiescent) post(si, .always, "{s}: ready", .{specs[si].name}); return; } } // ------------------------------------------------------------- introspection fn status(req: lsp.Req, out: *std.Io.Writer) !void { _ = req; try out.print("\nprotocol servers (lsp-client):\n", .{}); for (&specs, 0..) |*s, i| { const c = &conns[i]; c.lock(); const state = c.state; var rootbuf: [256]u8 = undefined; const rootlen = @min(c.root.len, rootbuf.len); @memcpy(rootbuf[0..rootlen], c.root[0..rootlen]); const root = rootbuf[0..rootlen]; const enc = c.caps.enc; const ndocs = c.docs.items.len; const ndiag = c.diags.items.len; const pull = c.caps.pull; const chier = c.caps.call_hier; const thier = c.caps.type_hier; c.unlock(); const env_v = if (getenv(s.env)) |v| std.mem.span(v) else null; try out.print(" {s:<28} {s}", .{ s.name, @tagName(state) }); if (state == .ready or state == .starting) { try out.print(" root {s} {s}", .{ root, @tagName(enc) }); if (pull) try out.print(" pull-diags", .{}); if (chier) try out.print(" call-hier", .{}); if (thier) try out.print(" type-hier", .{}); try out.print(" docs {d} diag-files {d}", .{ ndocs, ndiag }); } if (env_v) |v| try out.print(" [{s}={s}]", .{ s.env, if (v.len == 0) "(disabled)" else v }); try out.print("\n ", .{}); for (s.exts) |e| try out.print("{s} ", .{e}); try out.print("\n", .{}); } while (!log_mu.tryLock()) std.atomic.spinLoopHint(); defer log_mu.unlock(); try out.print("\nclient queries ({d} total, keeping {d}):\n", .{ log_total, log_cap }); var shown: usize = 0; for (0..log_cap) |i| { const e = &log_buf[(log_next + i) % log_cap]; if (!e.used) continue; shown += 1; if (hideTime()) { try out.print(" {s:<18} {s:<14} {s:<20} {d:>4} row(s){s}{s}\n", .{ @tagName(e.kind), specs[e.spec].name, e.file[0..e.file_len], e.rows, if (e.err_len > 0) " ERROR " else "", e.err[0..e.err_len], }); } else { try out.print(" {s:<18} {s:<14} {s:<20} {d:>7}us {d:>4} row(s){s}{s}\n", .{ @tagName(e.kind), specs[e.spec].name, e.file[0..e.file_len], e.us, e.rows, if (e.err_len > 0) " ERROR " else "", e.err[0..e.err_len], }); } } if (shown == 0) try out.print(" (none yet — press gd in a file with a server, then ask again)\n", .{}); } // ------------------------------------------------------------------ plumbing fn app(b: *std.ArrayList(u8), arena: std.mem.Allocator, s: []const u8) Err!void { b.appendSlice(arena, s) catch return error.OutOfMemory; } /// `"textDocument":{"uri":U},"position":{...}` — the params shared by every /// position request. fn atPos(arena: std.mem.Allocator, uri: []const u8, pos: Pos) Err!std.ArrayList(u8) { var b: std.ArrayList(u8) = .empty; try app(&b, arena, "\"textDocument\":{\"uri\":"); try jstr(&b, arena, uri); b.print(arena, "}},\"position\":{{\"line\":{d},\"character\":{d}}}", .{ pos.line, pos.ch }) catch return error.OutOfMemory; // the extra `}` above closed textDocument; nothing to fix up return b; } fn frame(sock: c_int, body: []const u8, deadline: i64) Err!void { if (sock < 0) return error.Dead; var hdr: [64]u8 = undefined; const h = std.fmt.bufPrint(&hdr, "Content-Length: {d}\r\n\r\n", .{body.len}) catch return error.Protocol; try sockWrite(sock, h, deadline); try sockWrite(sock, body, deadline); } fn sockWrite(sock: c_int, bytes: []const u8, deadline: i64) Err!void { var off: usize = 0; while (off < bytes.len) { const left = deadline - nowMs(); if (left <= 0) return error.Timeout; var pfd = [_]libc.pollfd{.{ .fd = sock, .events = POLLOUT, .revents = 0 }}; const pr = libc.poll(&pfd, 1, @intCast(@min(left, 1000))); if (pr < 0) { if (libc.errno(pr) == .INTR) continue; return error.Dead; } if (pr == 0) continue; const n = libc.send(sock, bytes.ptr + off, bytes.len - off, msg_nosignal); if (n < 0) { const e = libc.errno(n); if (e == .INTR or e == .AGAIN) continue; return error.Dead; } if (n == 0) return error.Dead; off += @intCast(n); } } fn nowMs() i64 { var ts: libc.timespec = undefined; _ = libc.clock_gettime(.MONOTONIC, &ts); return @as(i64, @intCast(ts.sec)) * 1000 + @divTrunc(@as(i64, @intCast(ts.nsec)), 1_000_000); } fn nowUs() u64 { var ts: libc.timespec = undefined; _ = libc.clock_gettime(.MONOTONIC, &ts); return @as(u64, @intCast(ts.sec)) *| 1_000_000 +| @as(u64, @intCast(ts.nsec)) / 1000; } /// The spec's binary: the env override when set, else `bin`, PATH-searched /// unless it names a path. fn binOf(spec: *const Spec, buf: *[4096:0]u8) ?[:0]const u8 { const name = if (getenv(spec.env)) |v| std.mem.span(v) else spec.bin; if (name.len == 0) return null; if (std.mem.indexOfScalar(u8, name, '/') != null) { const z = std.fmt.bufPrintSentinel(buf, "{s}", .{name}, 0) catch return null; return if (access(z.ptr, X_OK) == 0) z else null; } const path = if (getenv("PATH")) |p| std.mem.span(p) else "/usr/local/bin:/usr/bin:/bin"; var it = std.mem.tokenizeScalar(u8, path, ':'); while (it.next()) |dir| { const cand = std.fmt.bufPrintSentinel(buf, "{s}/{s}", .{ dir, name }, 0) catch continue; if (access(cand.ptr, X_OK) == 0) return cand; } return null; } /// Where the project starts: helix's find_root rule. Walk UP from `dir`; the /// TOP-MOST directory holding a language marker wins (a cargo workspace's /// root Cargo.toml beats the member crate's), the CLOSEST `.git` is the /// fallback, the asking directory the last resort. fn rootOf(dir: []const u8, markers: []const []const u8) []const u8 { if (dir.len == 0 or dir[0] != '/') return "/"; var top_marker: ?[]const u8 = null; var git: ?[]const u8 = null; var d = dir; var buf: [4096:0]u8 = undefined; while (true) { for (markers) |marker| { const p = std.fmt.bufPrintSentinel(&buf, "{s}/{s}", .{ d, marker }, 0) catch continue; if (access(p.ptr, 0) == 0) { top_marker = d; break; } } if (git == null) { if (std.fmt.bufPrintSentinel(&buf, "{s}/.git", .{d}, 0)) |p| { if (access(p.ptr, 0) == 0) git = d; } else |_| {} } d = std.fs.path.dirname(d) orelse break; if (d.len <= 1) break; } return top_marker orelse git orelse dir; } /// A path as a `file://` URI, raw (no JSON quotes). Everything outside the /// unreserved set is percent-encoded; `/` stays a separator. fn uriOf(b: *std.ArrayList(u8), arena: std.mem.Allocator, path: []const u8) Err!void { b.appendSlice(arena, "file://") catch return error.OutOfMemory; for (path) |ch| { if (std.ascii.isAlphanumeric(ch) or ch == '/' or ch == '-' or ch == '_' or ch == '.' or ch == '~') { b.append(arena, ch) catch return error.OutOfMemory; } else { b.print(arena, "%{X:0>2}", .{ch}) catch return error.OutOfMemory; } } } /// `file:///a/b%20c` -> `/a/b c`. Rows carry real paths; look.zig opens them /// and `%20` is not a filename. A decoded CONTROL byte rejects the whole uri: /// a `%0A` in a path would split one location into two physical rows — a /// server-forged extra row in the results buffer — and no path worth opening /// has a newline, tab or NUL in it. fn pathOf(arena: std.mem.Allocator, uri: []const u8) ?[]const u8 { var rest = uri; if (std.mem.startsWith(u8, rest, "file://")) { rest = rest["file://".len..]; if (rest.len > 0 and rest[0] != '/') return null; // an authority we cannot open } else if (std.mem.indexOf(u8, rest, "://") != null) return null; var out: std.ArrayList(u8) = .empty; var i: usize = 0; while (i < rest.len) { var byte: u8 = rest[i]; if (rest[i] == '%' and i + 2 < rest.len) { if (std.fmt.parseInt(u8, rest[i + 1 .. i + 3], 16)) |b| { byte = b; i += 3; } else |_| i += 1; } else i += 1; if (byte < 0x20) return null; out.append(arena, byte) catch return null; } return out.items; } /// A JSON string. Invalid UTF-8 becomes `?` one byte at a time rather than /// U+FFFD: a buffer being typed into may be invalid mid-keystroke, and a /// replacement that changed byte lengths would move every offset after it. fn jstr(b: *std.ArrayList(u8), arena: std.mem.Allocator, s: []const u8) Err!void { b.append(arena, '"') catch return error.OutOfMemory; var i: usize = 0; while (i < s.len) { const ch = s[i]; if (ch < 0x80) { switch (ch) { '"' => b.appendSlice(arena, "\\\"") catch return error.OutOfMemory, '\\' => b.appendSlice(arena, "\\\\") catch return error.OutOfMemory, '\n' => b.appendSlice(arena, "\\n") catch return error.OutOfMemory, '\r' => b.appendSlice(arena, "\\r") catch return error.OutOfMemory, '\t' => b.appendSlice(arena, "\\t") catch return error.OutOfMemory, else => if (ch < 0x20) b.print(arena, "\\u{x:0>4}", .{ch}) catch return error.OutOfMemory else b.append(arena, ch) catch return error.OutOfMemory, } i += 1; continue; } const l = std.unicode.utf8ByteSequenceLength(ch) catch { b.append(arena, '?') catch return error.OutOfMemory; i += 1; continue; }; if (i + l > s.len or !std.unicode.utf8ValidateSlice(s[i .. i + l])) { b.append(arena, '?') catch return error.OutOfMemory; i += 1; continue; } b.appendSlice(arena, s[i .. i + l]) catch return error.OutOfMemory; i += l; } b.append(arena, '"') catch return error.OutOfMemory; } const Pos = struct { line: u32, ch: u32 }; /// Byte offset -> LSP position. `lsp.lineCol` gives the line and the BYTE /// column; only the column needs re-counting, and only when the server /// refused utf-8. fn posOf(src: []const u8, off: u32, enc: Enc) Pos { const lc = lsp.lineCol(src, off); if (enc == .utf8) return .{ .line = @intCast(lc.line), .ch = @intCast(lc.col) }; const bol = @min(off, src.len) - lc.col; var units: u32 = 0; var i: usize = bol; while (i < bol + lc.col) { const l = std.unicode.utf8ByteSequenceLength(src[i]) catch 1; units += if (l == 4) 2 else 1; i += l; } return .{ .line = @intCast(lc.line), .ch = units }; } // std.json.Value navigation: optional-in, optional-out, so a missing field // and a wrongly-typed one read the same and neither can panic on a server // that sends the unexpected. fn get(v: ?std.json.Value, key: []const u8) ?std.json.Value { const o = v orelse return null; if (o != .object) return null; return o.object.get(key); } fn str(v: ?std.json.Value) ?[]const u8 { const o = v orelse return null; return if (o == .string) o.string else null; } fn num(v: ?std.json.Value) ?i64 { const o = v orelse return null; return if (o == .integer) o.integer else null; } fn items(v: ?std.json.Value) []std.json.Value { const o = v orelse return &.{}; return if (o == .array) o.array.items else &.{}; } const Range = struct { sl: u32, sc: u32, el: u32, ec: u32 }; fn rangeOf(range: ?std.json.Value) ?Range { const st = get(range, "start") orelse return null; const en = get(range, "end") orelse return null; return .{ .sl = coord(get(st, "line")) orelse return null, .sc = coord(get(st, "character")) orelse return null, .el = coord(get(en, "line")) orelse return null, .ec = coord(get(en, "character")) orelse return null, }; } /// One protocol coordinate, VALIDATED rather than clamped: a negative line is /// not "line zero", it is a malformed response, and a value past u32 would /// panic the @intCast that follows. Server bytes are not trusted bytes. fn coord(v: ?std.json.Value) ?u32 { const n = num(v) orelse return null; if (n < 0 or n > std.math.maxInt(u32)) return null; return @intCast(n); } // ----------------------------------------------------------------- tests test "LSP client rename and format propagate incomplete edit encoding" { const gpa = std.testing.allocator; const uri = "file:///file.c"; const parsed = try std.json.parseFromSlice(std.json.Value, gpa, \\{"changes":{"file:///file.c":[{"range":{"start":{"line":0,"character":0},"end":{"line":0,"character":3}},"newText":"A%\n"},{"range":{"start":{"line":0,"character":4},"end":{"line":0,"character":7}},"newText":"B"}]}} , .{}); defer parsed.deinit(); const expected = "@put 0 3 A%25%0A\n@put 4 7 B\n"; for ([_]bool{ true, false }) |rename| { var buffer: [128]u8 = undefined; for (0..expected.len + 1) |capacity| { var arena: std.heap.ArenaAllocator = .init(gpa); defer arena.deinit(); var out: std.Io.Writer = .fixed(buffer[0..capacity]); var cx: Cx = .{ .arena = arena.allocator(), .base = "/", .cur_path = "/file.c", .cur_src = "abc xyz", .out = &out }; const result = if (rename) renameEdits(&cx, .utf8, uri, parsed.value, cx.cur_src) else formatEdits(&cx, .utf8, get(get(parsed.value, "changes"), uri).?, cx.cur_src); if (capacity < expected.len) { try std.testing.expectError(error.WriteFailed, result); } else { try result; try std.testing.expectEqualStrings(expected, out.buffered()); } } } } test "frameNext distinguishes incomplete, valid and poison frames" { try std.testing.expectEqual(FrameStep.incomplete, frameNext("Content-Length: 5\r\n")); try std.testing.expectEqual(FrameStep.incomplete, frameNext("Content-Length: 5\r\n\r\nhel")); const one = frameNext("Content-Length: 5\r\n\r\nhello").frame; try std.testing.expectEqualStrings("hello", "Content-Length: 5\r\n\r\nhello"[one.body_start..one.body_end]); const two = "content-length: 2\r\nX-Other: y\r\n\r\nab" ++ "Content-Length: 1\r\n\r\nz"; const first = frameNext(two).frame; try std.testing.expectEqualStrings("ab", two[first.body_start..first.body_end]); const second = frameNext(two[first.consumed..]).frame; try std.testing.expectEqualStrings("z", two[first.consumed..][second.body_start..second.body_end]); // complete-but-unusable headers can never re-synchronize: poison, not // "keep buffering" — the wedge a review found and this line pins try std.testing.expectEqual(FrameStep.poison, frameNext("Content-Length: nope\r\n\r\n")); try std.testing.expectEqual(FrameStep.poison, frameNext("X-Only: y\r\n\r\n")); try std.testing.expectEqual(FrameStep.poison, frameNext("Content-Length: 67108865\r\n\r\n")); } test "server coordinates are validated, not clamped" { var arena_state: std.heap.ArenaAllocator = .init(std.testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); const bad = [_][]const u8{ \\{"start":{"line":-1,"character":0},"end":{"line":0,"character":1}} , \\{"start":{"line":0,"character":0},"end":{"line":0,"character":4294967296}} , \\{"start":{"line":0,"character":0},"end":{"line":0}} }; for (bad) |json| { const v = try std.json.parseFromSliceLeaky(std.json.Value, arena, json, .{}); try std.testing.expect(rangeOf(v) == null); } // strict offsets: past-EOL and past-EOF are refusals, EOL itself is real const src = "ab\ncd"; try std.testing.expectEqual(@as(?usize, 2), strictOffset(src, 0, 2, .utf8)); try std.testing.expect(strictOffset(src, 0, 3, .utf8) == null); try std.testing.expect(strictOffset(src, 2, 0, .utf8) == null); try std.testing.expectEqual(@as(?usize, 5), strictOffset(src, 1, 2, .utf8)); } test "uri round-trips spaces and utf-8" { var arena_state: std.heap.ArenaAllocator = .init(std.testing.allocator); defer arena_state.deinit(); const arena = arena_state.allocator(); var b: std.ArrayList(u8) = .empty; try uriOf(&b, arena, "/a dir/naïve.rs"); try std.testing.expectEqualStrings("file:///a%20dir/na%C3%AFve.rs", b.items); try std.testing.expectEqualStrings("/a dir/naïve.rs", pathOf(arena, b.items).?); try std.testing.expect(pathOf(arena, "file://host/x") == null); // a %0A would split one location into two rows: rejected wholesale try std.testing.expect(pathOf(arena, "file:///tmp/a%0A/tmp/b:1:1") == null); try std.testing.expect(pathOf(arena, "file:///tmp/a%00b") == null); try std.testing.expect(pathOf(arena, "https://x/y") == null); } test "posOf and offsetAt invert each other under both encodings" { const src = "aé𝕏b\ncd\n"; // é: 2 bytes/1 unit, 𝕏: 4 bytes/2 units inline for (.{ Enc.utf8, Enc.utf16 }) |enc| { const off: u32 = 7; // the 'b' const p = posOf(src, off, enc); try std.testing.expectEqual(@as(u32, 0), p.line); try std.testing.expectEqual(off, @as(u32, @intCast(offsetAt(src, p.line, p.ch, enc)))); } const p2 = posOf("aé𝕏b\ncd\n", 10, .utf16); // the 'd' try std.testing.expectEqual(@as(u32, 1), p2.line); try std.testing.expectEqual(@as(u32, 1), p2.ch); } test "rootOf takes the top-most marker and falls back to git then dir" { // libc mkdtemp under /tmp rather than std.testing.tmpDir: rootOf probes // with access(2) on ABSOLUTE paths, and /tmp has no ancestor markers to // muddy the fallback assertions the way the repo's own .zig-cache does. var tpl: [64:0]u8 = undefined; _ = std.fmt.bufPrintSentinel(&tpl, "/tmp/pardes-rootof-XXXXXX", .{}, 0) catch unreachable; const base_z = mkdtemp(&tpl) orelse return error.TestUnexpectedResult; const base = std.mem.span(base_z); defer { var cmd: [128:0]u8 = undefined; if (std.fmt.bufPrintSentinel(&cmd, "rm -rf {s}", .{base}, 0)) |z| _ = system(z.ptr) else |_| {} } var pb: [128:0]u8 = undefined; for ([_][]const u8{ "/ws", "/ws/member", "/ws/member/src" }) |d| { const z = try std.fmt.bufPrintSentinel(&pb, "{s}{s}", .{ base, d }, 0); try std.testing.expect(mkdir(z.ptr, 0o700) == 0); } for ([_][]const u8{ "/ws/Cargo.toml", "/ws/member/Cargo.toml" }) |f| { const z = try std.fmt.bufPrintSentinel(&pb, "{s}{s}", .{ base, f }, 0); const fd = libc.open(z.ptr, .{ .ACCMODE = .WRONLY, .CREAT = true }, @as(libc.mode_t, 0o600)); try std.testing.expect(fd >= 0); _ = libc.close(fd); } const deep = try std.fmt.bufPrintSentinel(&pb, "{s}/ws/member/src", .{base}, 0); const markers = [_][]const u8{"Cargo.toml"}; const root = rootOf(deep, &markers); try std.testing.expect(std.mem.endsWith(u8, root, "/ws")); // no language marker and no .git anywhere under /tmp: the asking dir wins const nothing = rootOf(deep, &.{"no-such-marker"}); try std.testing.expect(std.mem.startsWith(u8, deep, nothing)); } test "specFor routes extensions and honours the disable env" { try std.testing.expect(specFor("/x/main.rs") != null); try std.testing.expect(specFor("/x/main.zig") == null); try std.testing.expect(specFor("/x/README.md") == null); try std.testing.expectEqualStrings("rust-analyzer", specs[specFor("/x/main.rs").?].name); } test "the transport this host actually gives us is a pair, and both ends are close-on-exec" { // The spawn's FIRST fallible step, and for one release on macOS its last: // `SOCK.CLOEXEC` is spelled for darwin in zig's libc bindings and rejected // by darwin's socketpair(2), so this returned EPROTONOSUPPORT and no // language server was ever forked on that platform. Nothing above the // transport can notice — `ensure` reports the same `NoServer` a missing // binary does — so the check belongs here, on the real function, in a test // that runs on the host rather than on the linux target the snapshot // suite cross-compiles to. var sv: [2]libc.fd_t = undefined; try std.testing.expect(transportPair(&sv)); defer { _ = libc.close(sv[0]); _ = libc.close(sv[1]); } // close-on-exec on both ends, however the platform got there: the flag on // linux, fcntl on darwin. Without it every pty shell forked afterwards // inherits the server's socket, which is the bug host_io.zig fixed for the // pty master. const FD_CLOEXEC: c_int = 1; for (sv) |fd| try std.testing.expect(libc.fcntl(fd, libc.F.GETFD, @as(c_int, 0)) & FD_CLOEXEC != 0); // and it is a connected PAIR, not two unrelated descriptors const msg = "ping"; try std.testing.expectEqual(@as(isize, msg.len), libc.write(sv[0], msg, msg.len)); var got: [8]u8 = undefined; try std.testing.expectEqual(@as(isize, msg.len), libc.read(sv[1], &got, got.len)); try std.testing.expectEqualStrings(msg, got[0..msg.len]); }