summaryrefslogtreecommitdiff
path: root/src/lsp/lsp_client.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-01 09:23:53 -0300
committerGabriel Schneider <[email protected]>2026-09-01 11:24:12 -0300
commitae9325a5cb128d0d952afb8f9feaaca68e5e37a2 (patch)
tree9ae44ac38f7b71edfe2882d0a882dbc61304ec80 /src/lsp/lsp_client.zig
parent848ad99fa597387a85f75e752dc4c9e10f8c24f4 (diff)
downloadpardes-ae9325a5cb128d0d952afb8f9feaaca68e5e37a2.tar.gz
pardes-ae9325a5cb128d0d952afb8f9feaaca68e5e37a2.zip
lsp: a protocol client for every other language, narrated on the message row
The seam grows a second backend: src/lsp/lsp_client.zig speaks JSON-RPC to child language servers — rust-analyzer, clangd, gopls, tsserver, pyright are rows in a spec table — while the in-process ZLS analyser keeps .zig. One reader thread per server owns the socket, routes responses to a mailbox under the conn mutex (monotonic condvar), answers server-to-client requests, feeds the diagnostics store, and narrates $/progress and state changes through a status sink both native shells post to the transient message row: "rust-analyzer: cargo check 88% 955/1083" lands where a save narrates, with the same clock. Chatty progress is throttled and deduplicated; settled states always land, which is also what makes the goldens deterministic. Nothing wedges and nothing healthy dies: waits are deadline-bounded, a timeout cancels and returns no rows, three consecutive timeouts restart the server ONLY while it is idle (an indexing server is narrating its own excuse), spawn and handshake failures back off 10s to 2min, a crash shortly after ready counts as a failure, and only a missing binary disables a spec. PARDES_LSP_{RS,C,GO,TS,PY} override binaries; empty disables; the snapshot harness pins RS to test/lspmock.zig and empties the rest. Mutating answers really mutate now: the @put record beside rename @edit carries per-range text, so = applies the formatter (both backends) and a same-file WorkspaceEdit rename applies atomically, one undo step, narrated ("renamed 2 range(s)"); a multi-file rename previews as rows instead of half-applying. Malformed responses fail closed: coordinates validated not clamped, one bad TextEdit poisons the whole edit set, poison frames kill the connection instead of buffering forever, decoded control bytes reject a uri, hierarchy items too deep to reserialize are skipped. Four kinds helix does not have, on SPC l: c/C incoming/outgoing calls (rows are call sites), t/T super/subtypes. Pull diagnostics (3.17) preferred when advertised. Help gains a language-keys footer for the motions no builtin row could carry; lsp.rel and look.grep now share one path-shortening rule. zig build lspprobe drives the seam from the CLI (comma-separated kinds share one server); measured against a 1083-crate workspace warm: gd 26ms, gr 213 rows 165ms, incoming calls 212 sites 197ms, document symbols 670 rows 347ms. docs/lsp.md tells the whole story; lsp-evaluation.md gets an addendum.
Diffstat (limited to 'src/lsp/lsp_client.zig')
-rw-r--r--src/lsp/lsp_client.zig2127
1 files changed, 2127 insertions, 0 deletions
diff --git a/src/lsp/lsp_client.zig b/src/lsp/lsp_client.zig
new file mode 100644
index 00000000..eca2691d
--- /dev/null
+++ b/src/lsp/lsp_client.zig
@@ -0,0 +1,2127 @@
+//! 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 };
+
+/// 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
+
+/// The seam entry point. Never fails, never panics; no rows is the only error
+/// rendering there is (`SPC l i` shows what was swallowed).
+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) catch {};
+
+ 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 = arena.alloc(u8, max_rows * 512) catch return;
+ var scratch: std.Io.Writer = .fixed(scratch_buf);
+
+ const t0 = nowUs();
+ var err_name: []const u8 = "";
+ answer(arena, si, req, &scratch, &tr) catch |e| {
+ err_name = @errorName(e);
+ tr.note("ERROR: {s} — the editor shows this as 'no result'", .{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);
+ out.writeAll(scratch.buffered()) catch {};
+}
+
+fn traceOut(tr: *const Trace, req: lsp.Req, out: *std.Io.Writer, rows: usize, us: u64) void {
+ if (req.kind != .explain) return;
+ 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),
+ }) catch {};
+ out.writeAll(tr.buf[0..tr.len]) catch {};
+ if (hideTime())
+ out.print("\n{d} row(s)\n", .{rows}) catch {}
+ else
+ out.print("\n{d} row(s) in {d}us\n", .{ rows, us }) catch {};
+}
+
+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 => {},
+ }
+ 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| 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;
+ out.writeAll(std.mem.trim(u8, text, " \t\r\n")) catch {};
+ out.writeByte('\n') catch {};
+ },
+ .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;
+ 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;
+ 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;
+ 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.
+ 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) {
+ 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");
+ emitRange(cx, enc, str(tu) orelse continue, r, "");
+ } else {
+ 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) void {
+ if (depth > 8) return;
+ for (items(node)) |sym| {
+ const name = str(get(sym, "name")) orelse continue;
+ if (get(get(sym, "location"), "range")) |r| {
+ emitRange(cx, enc, str(get(get(sym, "location"), "uri")) orelse uri, r, name);
+ } else {
+ emitRange(cx, enc, uri, get(sym, "selectionRange") orelse get(sym, "range"), name);
+ }
+ if (get(sym, "children")) |kids| 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| 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);
+ 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| emitDiag(cx, c.caps.enc, uri, dg, arena);
+ }
+ return;
+ }
+ 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) 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 = std.fmt.allocPrint(arena, "{s}: {s}", .{ label, flat(arena, str(get(dg, "message")) orelse "") }) catch return;
+ emitRange(cx, enc, uri, get(dg, "range"), msg);
+}
+
+fn renderStore(c: *Conn, arena: std.mem.Allocator, cx: *Cx, only_uri: ?[]const u8) 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| 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 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 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);
+ 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| 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| 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) {
+ emitRange(cx, c.caps.enc, fu, get(from, "selectionRange"), name);
+ } else for (ranges) |r| emitRange(cx, c.caps.enc, fu, r, name);
+ },
+ .outgoing => {
+ const to = get(entry, "to") orelse continue;
+ emitRange(cx, c.caps.enc, str(get(to, "uri")) orelse continue, get(to, "selectionRange") orelse get(to, "range"), hierText(cx.arena, to));
+ },
+ .supers, .subs => 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) 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);
+ lsp.spanRow(cx.out, shown, r.sl, col, r.el, end_col, rowtext);
+ } else {
+ 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, " ");
+}
+
+// ----------------------------------------------------------- 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 (libc.socketpair(libc.AF.UNIX, libc.SOCK.STREAM | libc.SOCK.CLOEXEC, 0, &sv) != 0) {
+ sa.free(owned_root);
+ return error.NoServer;
+ }
+ if (comptime builtin.os.tag.isDarwin()) {
+ const one: c_int = 1;
+ _ = libc.setsockopt(sv[0], libc.SOL.SOCKET, so_nosigpipe, @ptrCast(&one), @sizeOf(c_int));
+ }
+ const pid = libc.fork();
+ if (pid < 0) {
+ _ = libc.close(sv[0]);
+ _ = libc.close(sv[1]);
+ sa.free(owned_root);
+ 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 "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);
+}