tornado rejects fetch(request, **kwargs) when the first arg is already an HTTPRequest and the kwargs overlap request construction (auth_username / auth_password). Move auth onto the HTTPRequest. Add OpenCodeClient.close() to release its AsyncHTTPClient; module-level atexit closes the tornado singleton so 'jupyter lab' can exit cleanly after SIGINT instead of hanging on 'received signal 2, stopping'. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
142 lines
4.8 KiB
Python
142 lines
4.8 KiB
Python
"""Async HTTP client for the local OpenCode server."""
|
|
|
|
import atexit
|
|
import json
|
|
from typing import Any, Optional
|
|
|
|
from tornado.httpclient import AsyncHTTPClient, HTTPRequest
|
|
from tornado.httputil import HTTPHeaders
|
|
|
|
from .config import OpenCodeConfig
|
|
|
|
|
|
class OpenCodeError(Exception):
|
|
"""Raised when an OpenCode API request fails."""
|
|
|
|
|
|
# Ensure the tornado AsyncHTTPClient singleton is closed at interpreter
|
|
# shutdown so the IOLoop can finish and the process can exit cleanly
|
|
# (otherwise Ctrl+C on `jupyter lab` hangs after "received signal 2,
|
|
# stopping" because the singleton keeps persistent connections alive).
|
|
def _close_opencode_http() -> None:
|
|
try:
|
|
AsyncHTTPClient().close()
|
|
except Exception:
|
|
# Never block interpreter shutdown on close errors.
|
|
pass
|
|
|
|
|
|
atexit.register(_close_opencode_http)
|
|
|
|
|
|
class OpenCodeClient:
|
|
def __init__(
|
|
self,
|
|
config: OpenCodeConfig,
|
|
http_client: Optional[AsyncHTTPClient] = None,
|
|
) -> None:
|
|
self._config = config
|
|
# Use a dedicated client (not the singleton) so close() can
|
|
# deterministically release its connections without affecting
|
|
# other components. The module-level atexit still closes the
|
|
# singleton as a safety net.
|
|
self._http_client = http_client or AsyncHTTPClient(force_instance=True)
|
|
self._owns_http_client = http_client is None
|
|
|
|
def close(self) -> None:
|
|
"""Release the underlying HTTP client's connections. Idempotent."""
|
|
if getattr(self, "_http_client", None) is None:
|
|
return
|
|
try:
|
|
if self._owns_http_client:
|
|
self._http_client.close()
|
|
except Exception:
|
|
pass
|
|
self._http_client = None # type: ignore[assignment]
|
|
self._owns_http_client = False
|
|
|
|
async def health(self) -> dict[str, Any]:
|
|
return await self._request("GET", "/global/health")
|
|
|
|
async def list_providers(self) -> list[dict[str, Any]]:
|
|
return await self._request("GET", "/config/providers")
|
|
|
|
async def create_session(self, title: str) -> dict[str, Any]:
|
|
return await self._request("POST", "/session", {"title": title})
|
|
|
|
async def send_message_sync(
|
|
self,
|
|
session_id: str,
|
|
parts: list[dict[str, Any]],
|
|
provider_id: Optional[str] = None,
|
|
model_id: Optional[str] = None,
|
|
system: Optional[str] = None,
|
|
) -> dict[str, Any]:
|
|
body: dict[str, Any] = {"parts": parts}
|
|
if provider_id is not None and model_id is not None:
|
|
body["model"] = {"providerID": provider_id, "modelID": model_id}
|
|
if system is not None:
|
|
body["system"] = system
|
|
return await self._request(
|
|
"POST",
|
|
"/session/%s/message" % session_id,
|
|
body,
|
|
)
|
|
|
|
async def abort(self, session_id: str) -> bool:
|
|
result = await self._request("POST", "/session/%s/abort" % session_id)
|
|
return result is not None
|
|
|
|
async def delete_session(self, session_id: str) -> bool:
|
|
result = await self._request("DELETE", "/session/%s" % session_id)
|
|
return result is not None
|
|
|
|
@property
|
|
def endpoint(self) -> str:
|
|
return self._config.url
|
|
|
|
async def _request(
|
|
self,
|
|
method: str,
|
|
path: str,
|
|
body: Optional[dict[str, Any]] = None,
|
|
) -> Optional[Any]:
|
|
url = self._config.url + path
|
|
headers = HTTPHeaders(
|
|
{
|
|
"Content-Type": "application/json",
|
|
"Accept": "application/json",
|
|
}
|
|
)
|
|
request_kwargs: dict[str, Any] = {
|
|
"method": method,
|
|
"headers": headers,
|
|
"request_timeout": self._config.request_timeout_seconds,
|
|
}
|
|
if body is not None:
|
|
request_kwargs["body"] = json.dumps(body).encode("utf-8")
|
|
# HTTP Basic Auth must be set on the HTTPRequest itself; passing
|
|
# auth_* as **kwargs to fetch() when the first arg is already an
|
|
# HTTPRequest is rejected by tornado ("kwargs can't be used if
|
|
# request is an HTTPRequest object").
|
|
if self._config.auth is not None:
|
|
request_kwargs["auth_username"] = self._config.auth[0]
|
|
request_kwargs["auth_password"] = self._config.auth[1]
|
|
|
|
request = HTTPRequest(url, **request_kwargs)
|
|
|
|
response = await self._http_client.fetch(request)
|
|
|
|
if 200 <= response.code < 300:
|
|
if not response.body:
|
|
return True
|
|
return json.loads(response.body.decode("utf-8"))
|
|
if response.code == 404:
|
|
return None
|
|
|
|
response_body = response.body.decode("utf-8") if response.body else ""
|
|
raise OpenCodeError(
|
|
"OpenCode request %s %s failed with status %s: %s"
|
|
% (method, url, response.code, response_body)
|
|
)
|