|
| 1 | +"""Small, dependency-free LSP/JSON-RPC client used by the agent harness. |
| 2 | +
|
| 3 | +This intentionally implements only the transport/lifecycle pieces needed by the |
| 4 | +agent-facing LSP tool. It is synchronous at the public boundary because the |
| 5 | +harness' Tool API is synchronous; a reader thread keeps server messages flowing. |
| 6 | +""" |
| 7 | + |
| 8 | +from __future__ import annotations |
| 9 | + |
| 10 | +import contextlib |
| 11 | +import json |
| 12 | +import os |
| 13 | +import queue |
| 14 | +import subprocess |
| 15 | +import threading |
| 16 | +from pathlib import Path |
| 17 | +from typing import Any |
| 18 | + |
| 19 | + |
| 20 | +class LSPError(RuntimeError): |
| 21 | + pass |
| 22 | + |
| 23 | + |
| 24 | +class LSPClient: |
| 25 | + def __init__(self, command: list[str], root: str, language_id: str, timeout: float = 30.0): |
| 26 | + self.command = command |
| 27 | + self.root = os.path.realpath(root) |
| 28 | + self.language_id = language_id |
| 29 | + self.timeout = timeout |
| 30 | + self.proc: subprocess.Popen[bytes] | None = None |
| 31 | + self._reader: threading.Thread | None = None |
| 32 | + self._write_lock = threading.Lock() |
| 33 | + self._state_lock = threading.Lock() |
| 34 | + self._next_id = 1 |
| 35 | + self._pending: dict[int, queue.Queue] = {} |
| 36 | + self._opened: dict[str, int] = {} |
| 37 | + self._texts: dict[str, str] = {} |
| 38 | + self._diagnostics: dict[str, list[dict[str, Any]]] = {} |
| 39 | + self._capabilities: dict[str, Any] = {} |
| 40 | + self.position_encoding = "utf-16" |
| 41 | + self._closed = False |
| 42 | + |
| 43 | + @property |
| 44 | + def alive(self) -> bool: |
| 45 | + return self.proc is not None and self.proc.poll() is None and not self._closed |
| 46 | + |
| 47 | + def start(self) -> None: |
| 48 | + if self.alive: |
| 49 | + return |
| 50 | + self._closed = False |
| 51 | + try: |
| 52 | + self.proc = subprocess.Popen( |
| 53 | + self.command, |
| 54 | + stdin=subprocess.PIPE, |
| 55 | + stdout=subprocess.PIPE, |
| 56 | + stderr=subprocess.DEVNULL, |
| 57 | + cwd=self.root, |
| 58 | + bufsize=0, |
| 59 | + ) |
| 60 | + except OSError as e: |
| 61 | + raise LSPError(f"failed to start LSP server {' '.join(self.command)!r}: {e}") from e |
| 62 | + self._reader = threading.Thread( |
| 63 | + target=self._read_loop, daemon=True, name=f"lsp-reader:{self.command[0]}" |
| 64 | + ) |
| 65 | + self._reader.start() |
| 66 | + try: |
| 67 | + result = self.request( |
| 68 | + "initialize", |
| 69 | + { |
| 70 | + "processId": os.getpid(), |
| 71 | + "rootUri": Path(self.root).as_uri(), |
| 72 | + "workspaceFolders": [ |
| 73 | + {"uri": Path(self.root).as_uri(), "name": Path(self.root).name} |
| 74 | + ], |
| 75 | + "capabilities": { |
| 76 | + "general": {"positionEncodings": ["utf-16", "utf-8"]}, |
| 77 | + "workspace": { |
| 78 | + "workspaceFolders": True, |
| 79 | + "symbol": {"dynamicRegistration": False}, |
| 80 | + }, |
| 81 | + "textDocument": { |
| 82 | + "definition": {"linkSupport": True}, |
| 83 | + "references": {}, |
| 84 | + "hover": {"contentFormat": ["markdown", "plaintext"]}, |
| 85 | + "documentSymbol": {"hierarchicalDocumentSymbolSupport": True}, |
| 86 | + "implementation": {"linkSupport": True}, |
| 87 | + "callHierarchy": {"dynamicRegistration": False}, |
| 88 | + }, |
| 89 | + }, |
| 90 | + "trace": "off", |
| 91 | + }, |
| 92 | + ) |
| 93 | + self._capabilities = result.get("capabilities", {}) if isinstance(result, dict) else {} |
| 94 | + encoding = self._capabilities.get("positionEncoding") |
| 95 | + if encoding in {"utf-8", "utf-16", "utf-32"}: |
| 96 | + self.position_encoding = encoding |
| 97 | + self.notify("initialized", {}) |
| 98 | + except Exception: |
| 99 | + self.close() |
| 100 | + raise |
| 101 | + |
| 102 | + def request(self, method: str, params: Any = None) -> Any: |
| 103 | + self.start_if_needed() |
| 104 | + with self._state_lock: |
| 105 | + request_id = self._next_id |
| 106 | + self._next_id += 1 |
| 107 | + waiter: queue.Queue = queue.Queue(maxsize=1) |
| 108 | + self._pending[request_id] = waiter |
| 109 | + message: dict[str, Any] = {"jsonrpc": "2.0", "id": request_id, "method": method} |
| 110 | + if params is not None: |
| 111 | + message["params"] = params |
| 112 | + try: |
| 113 | + self._send(message) |
| 114 | + kind, value = waiter.get(timeout=self.timeout) |
| 115 | + except queue.Empty as e: |
| 116 | + with self._state_lock: |
| 117 | + self._pending.pop(request_id, None) |
| 118 | + raise LSPError(f"LSP request timed out: {method}") from e |
| 119 | + if kind == "error": |
| 120 | + code = value.get("code") if isinstance(value, dict) else None |
| 121 | + msg = ( |
| 122 | + value.get("message", "unknown LSP error") if isinstance(value, dict) else str(value) |
| 123 | + ) |
| 124 | + raise LSPError(f"LSP {method} failed ({code}): {msg}") |
| 125 | + return value |
| 126 | + |
| 127 | + def notify(self, method: str, params: Any = None) -> None: |
| 128 | + self.start_if_needed() |
| 129 | + message: dict[str, Any] = {"jsonrpc": "2.0", "method": method} |
| 130 | + if params is not None: |
| 131 | + message["params"] = params |
| 132 | + self._send(message) |
| 133 | + |
| 134 | + def start_if_needed(self) -> None: |
| 135 | + if not self.alive: |
| 136 | + self.start() |
| 137 | + |
| 138 | + def open_document(self, uri: str, text: str) -> None: |
| 139 | + self.start_if_needed() |
| 140 | + with self._state_lock: |
| 141 | + version = self._opened.get(uri, 0) |
| 142 | + if version: |
| 143 | + if self._texts.get(uri) == text: |
| 144 | + return |
| 145 | + version += 1 |
| 146 | + self._opened[uri] = version |
| 147 | + self._texts[uri] = text |
| 148 | + self.notify( |
| 149 | + "textDocument/didChange", |
| 150 | + { |
| 151 | + "textDocument": {"uri": uri, "version": version}, |
| 152 | + "contentChanges": [{"text": text}], |
| 153 | + }, |
| 154 | + ) |
| 155 | + return |
| 156 | + self._opened[uri] = 1 |
| 157 | + self._texts[uri] = text |
| 158 | + self.notify( |
| 159 | + "textDocument/didOpen", |
| 160 | + { |
| 161 | + "textDocument": { |
| 162 | + "uri": uri, |
| 163 | + "languageId": self.language_id, |
| 164 | + "version": 1, |
| 165 | + "text": text, |
| 166 | + } |
| 167 | + }, |
| 168 | + ) |
| 169 | + |
| 170 | + def close_document(self, uri: str) -> None: |
| 171 | + with self._state_lock: |
| 172 | + if uri not in self._opened: |
| 173 | + return |
| 174 | + self._opened.pop(uri, None) |
| 175 | + self._texts.pop(uri, None) |
| 176 | + if self.alive: |
| 177 | + self.notify("textDocument/didClose", {"textDocument": {"uri": uri}}) |
| 178 | + |
| 179 | + def diagnostics(self, uri: str) -> list[dict[str, Any]]: |
| 180 | + with self._state_lock: |
| 181 | + return list(self._diagnostics.get(uri, [])) |
| 182 | + |
| 183 | + def close(self) -> None: |
| 184 | + if self._closed: |
| 185 | + return |
| 186 | + proc = self.proc |
| 187 | + was_alive = proc is not None and proc.poll() is None |
| 188 | + if was_alive: |
| 189 | + with contextlib.suppress(Exception): |
| 190 | + self.request("shutdown", None) |
| 191 | + with contextlib.suppress(Exception): |
| 192 | + self.notify("exit", None) |
| 193 | + self._closed = True |
| 194 | + proc = self.proc |
| 195 | + if proc is not None and proc.poll() is None: |
| 196 | + with contextlib.suppress(Exception): |
| 197 | + proc.terminate() |
| 198 | + proc.wait(timeout=2) |
| 199 | + with contextlib.suppress(Exception): |
| 200 | + proc.kill() |
| 201 | + with self._state_lock: |
| 202 | + with contextlib.suppress(queue.Full): |
| 203 | + for waiter in self._pending.values(): |
| 204 | + waiter.put_nowait(("error", {"message": "LSP client closed"})) |
| 205 | + self._pending.clear() |
| 206 | + self._opened.clear() |
| 207 | + self._texts.clear() |
| 208 | + |
| 209 | + def _send(self, message: dict[str, Any]) -> None: |
| 210 | + data = json.dumps(message, ensure_ascii=False, separators=(",", ":")).encode("utf-8") |
| 211 | + packet = b"Content-Length: " + str(len(data)).encode("ascii") + b"\r\n\r\n" + data |
| 212 | + proc = self.proc |
| 213 | + if proc is None or proc.stdin is None or proc.poll() is not None: |
| 214 | + raise LSPError("LSP server is not running") |
| 215 | + with self._write_lock: |
| 216 | + try: |
| 217 | + proc.stdin.write(packet) |
| 218 | + proc.stdin.flush() |
| 219 | + except OSError as e: |
| 220 | + raise LSPError(f"failed writing to LSP server: {e}") from e |
| 221 | + |
| 222 | + def _read_loop(self) -> None: |
| 223 | + proc = self.proc |
| 224 | + if proc is None or proc.stdout is None: |
| 225 | + return |
| 226 | + stream = proc.stdout |
| 227 | + try: |
| 228 | + while not self._closed: |
| 229 | + headers: dict[str, str] = {} |
| 230 | + while True: |
| 231 | + line = stream.readline() |
| 232 | + if not line: |
| 233 | + return |
| 234 | + line = line.rstrip(b"\r\n") |
| 235 | + if not line: |
| 236 | + break |
| 237 | + if b":" in line: |
| 238 | + key, value = line.split(b":", 1) |
| 239 | + headers[key.decode("ascii", "replace").lower()] = value.strip().decode( |
| 240 | + "ascii", "replace" |
| 241 | + ) |
| 242 | + try: |
| 243 | + length = int(headers.get("content-length", "0")) |
| 244 | + except ValueError: |
| 245 | + continue |
| 246 | + payload = stream.read(length) |
| 247 | + if len(payload) != length: |
| 248 | + return |
| 249 | + try: |
| 250 | + message = json.loads(payload.decode("utf-8")) |
| 251 | + except (UnicodeDecodeError, json.JSONDecodeError): |
| 252 | + continue |
| 253 | + self._handle_message(message) |
| 254 | + finally: |
| 255 | + with self._state_lock: |
| 256 | + pending = list(self._pending.values()) |
| 257 | + self._pending.clear() |
| 258 | + with contextlib.suppress(queue.Full): |
| 259 | + for waiter in pending: |
| 260 | + waiter.put_nowait(("error", {"message": "LSP server exited"})) |
| 261 | + |
| 262 | + def _handle_message(self, message: dict[str, Any]) -> None: |
| 263 | + if "id" in message and ("result" in message or "error" in message): |
| 264 | + request_id = message.get("id") |
| 265 | + if isinstance(request_id, int): |
| 266 | + with self._state_lock: |
| 267 | + waiter = self._pending.pop(request_id, None) |
| 268 | + if waiter is not None: |
| 269 | + waiter.put( |
| 270 | + ("error", message["error"]) |
| 271 | + if "error" in message |
| 272 | + else ("result", message.get("result")) |
| 273 | + ) |
| 274 | + return |
| 275 | + |
| 276 | + method = message.get("method") |
| 277 | + params = message.get("params") |
| 278 | + if method == "textDocument/publishDiagnostics" and isinstance(params, dict): |
| 279 | + uri = params.get("uri") |
| 280 | + if isinstance(uri, str): |
| 281 | + with self._state_lock: |
| 282 | + self._diagnostics[uri] = params.get("diagnostics", []) or [] |
| 283 | + return |
| 284 | + |
| 285 | + # Servers may send requests/notifications to the client. Respond to |
| 286 | + # common client-side requests so a server does not stall waiting for |
| 287 | + # configuration/progress/UI handling that the coding agent does not need. |
| 288 | + if "id" in message and isinstance(message.get("id"), int): |
| 289 | + method = str(method or "") |
| 290 | + if method == "workspace/configuration": |
| 291 | + items = params if isinstance(params, list) else [] |
| 292 | + result = [None for _ in items] |
| 293 | + elif method == "workspace/applyEdit": |
| 294 | + result = {"applied": False} |
| 295 | + else: |
| 296 | + result = None |
| 297 | + with contextlib.suppress(LSPError): |
| 298 | + self._send({"jsonrpc": "2.0", "id": message["id"], "result": result}) |
0 commit comments