Files
notebook-ai-extension/opencode_bridge/opencode_client.py
T
b7089519fd feat: async + SSE message flow with interactive permission/question UI
Replace the synchronous /edit (wait for full markdown response) with an
async + SSE flow per demo.html so the user sees text stream in real-time
and can interact with permission/question events the agent raises.

Backend
-------
- EditHandler now calls OpenCode /session/:id/prompt_async and returns
  immediately with {ok, sessionId, notebookPath}; the LLM reply is no
  longer embedded in this response.
- New GlobalEventHandler proxies OpenCode /global/event as
  text/event-stream. Server forwards ALL events; the client filters.
  A too-eager server-side ?session= filter was silently dropping events
  the client would have accepted, so it was removed.
- New PermissionReplyHandler + QuestionReplyHandler forward user
  replies (once/always/reject and freeform answer) back to OpenCode
  Serve at /session/:sid/permissions/:permId and
  /session/:sid/question/:qId/reply.
- OpenCodeClient gains send_message_async(), stream_global_events()
  (async generator over the SSE feed via tornado streaming_callback),
  reply_permission() and reply_question(). Legacy send_message_sync
  removed.

Frontend
--------
- subscribeOpenCodeEvents() opens a fetch+reader SSE client (XSRF
  token injected from serverSettings); AbortController-backed close()
  is idempotent.
- OpenCodeInlinePrompt.applyEvent() routes events into the streaming
  UI: text delta -> assistant message (re-rendered as markdown on
  every delta so the user sees formatted <pre><code> blocks in
  real-time, not raw fence source); reasoning/tool/permission/question
  get collapsible details blocks. session.idle resets stream pointers
  and fires onStreamEnd.
- Permission/question blocks render real interactive UI: three
  buttons (once/always/reject) for permission, a text input + submit
  for question. Click handlers post through the new API routes and
  show success/failure status in place.
- Centralized idle detection (4 shapes: top-level + payload-nested
  x {session.idle, session.status/idle}) inside the prompt; cell
  action reacts via onStreamEnd callback rather than re-parsing
  event types.
- System / workspace / pty / lsp / mcp / installation events AND
  session-level control events (agent.switched, model.switched,
  file.edited) are not rendered in the frontend.

Tests
-----
- 52 backend pytest (was 43)
- 58 frontend jest (was 46)

design.md section 3 updated for the new async /edit response shape
and the new GET /events SSE endpoint.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-07-27 17:38:53 +08:00

308 lines
11 KiB
Python

"""Async HTTP client for the local OpenCode server.
Two interaction modes:
- **Request/response** (one-shot) - ``health``, ``list_providers``,
``create_session``, ``send_message_async``, ``abort``, ``delete_session``,
``list_session_messages``. ``send_message_async`` fires the prompt and
returns immediately; the actual LLM response is consumed via the SSE
stream (next mode).
- **Streaming** (long-lived) - :meth:`stream_global_events` is an async
generator over OpenCode Serve's ``GET /global/event`` SSE feed. It
yields each event as a parsed JSON dict until the upstream closes.
"""
import asyncio
import atexit
import json
import logging
from typing import Any, AsyncIterator, Optional
from tornado.httpclient import AsyncHTTPClient, HTTPRequest
from tornado.httputil import HTTPHeaders
from .config import OpenCodeConfig
log = logging.getLogger("opencode_bridge.opencode_client")
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_async(
self,
session_id: str,
parts: list[dict[str, Any]],
provider_id: Optional[str] = None,
model_id: Optional[str] = None,
system: Optional[str] = None,
) -> None:
"""Fire-and-forget POST /session/:id/prompt_async.
Returns once OpenCode has accepted the prompt and queued it for
processing - NOT when the LLM is done. Subscribe to
:meth:`stream_global_events` to follow the agent's progress.
"""
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
await self._request(
"POST",
"/session/%s/prompt_async" % 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 reply_permission(
self,
session_id: str,
permission_id: str,
response: str,
) -> bool:
"""Respond to a `permission.asked` event.
``response`` must be one of ``"once"`` / ``"always"`` / ``"reject"``,
matching the OpenCode Serve API. Returns True on success.
"""
result = await self._request(
"POST",
"/session/%s/permissions/%s" % (session_id, permission_id),
{"response": response},
)
return result is not None
async def reply_question(
self,
session_id: str,
question_id: str,
answer: str,
) -> bool:
"""Reply to a `question.asked` event with the user's freeform answer.
Returns True on success.
"""
result = await self._request(
"POST",
"/session/%s/question/%s/reply" % (session_id, question_id),
{"answer": answer},
)
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
async def list_session_messages(self, session_id: str) -> list[dict[str, Any]]:
"""List all messages in the given OpenCode session.
Returns the raw OpenCode response: a list of
``{ info: { role: "user"|"assistant", ... }, parts: [...] }``.
The server extension is responsible for projecting this into a
frontend-friendly ``{role, content}[]`` shape.
"""
return await self._request("GET", "/session/%s/message" % session_id)
@property
def endpoint(self) -> str:
return self._config.url
async def stream_global_events(self) -> AsyncIterator[dict[str, Any]]:
"""Stream OpenCode Serve's ``GET /global/event`` SSE feed.
Yields one parsed JSON dict per ``data: {...}`` event block. The
stream runs until:
- the upstream disconnects (end of stream),
- the consumer stops iterating (caller cancels this generator),
- or a network error occurs (logged and the generator returns).
The consumer is responsible for filtering events by session ID
(the upstream pushes events for ALL sessions, plus system events).
"""
queue: asyncio.Queue = asyncio.Queue()
end_sentinel = object()
def on_chunk(chunk: bytes) -> None:
if chunk:
queue.put_nowait(chunk)
headers = HTTPHeaders({
"Accept": "text/event-stream",
"Cache-Control": "no-cache",
})
request_kwargs: dict[str, Any] = {
"method": "GET",
"headers": headers,
"request_timeout": 0, # 0 = no timeout for long-lived streams
# streaming_callback must live on the HTTPRequest itself; tornado
# rejects any fetch() kwargs when an HTTPRequest is passed.
"streaming_callback": on_chunk,
}
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(self._config.url + "/global/event", **request_kwargs)
async def _run_fetch() -> None:
try:
await self._http_client.fetch(request, raise_error=False)
except Exception as e:
log.warning("OpenCode /global/event fetch error: %s", e)
finally:
queue.put_nowait(end_sentinel)
runner_task = asyncio.create_task(_run_fetch())
buffer = ""
try:
while True:
chunk = await queue.get()
if chunk is end_sentinel:
break
if not isinstance(chunk, bytes):
continue
buffer += chunk.decode("utf-8", errors="replace")
while "\n\n" in buffer:
block, buffer = buffer.split("\n\n", 1)
parsed = self._parse_sse_block(block)
if parsed is not None:
yield parsed
# Drain any trailing data after the last \n\n
if buffer.strip():
parsed = self._parse_sse_block(buffer)
if parsed is not None:
yield parsed
finally:
if not runner_task.done():
runner_task.cancel()
try:
await runner_task
except (asyncio.CancelledError, Exception):
pass
@staticmethod
def _parse_sse_block(block: str) -> Optional[dict[str, Any]]:
"""Parse a single SSE event block into a JSON dict (or None).
SSE format reminder: a block is one or more ``field: value`` lines
separated by newlines; multiple ``data:`` lines are concatenated
with ``\\n``. We only care about ``data:``.
"""
data_parts: list[str] = []
for line in block.split("\n"):
if line.startswith("data:"):
data_parts.append(line[5:].lstrip())
if not data_parts:
return None
data = "\n".join(data_parts).strip()
if not data:
return None
try:
result = json.loads(data)
except json.JSONDecodeError:
return None
return result if isinstance(result, dict) else None
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)
)