diff --git a/CHANGELOG.md b/CHANGELOG.md
index f923ea6..8734bb4 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -6,6 +6,14 @@ The format follows [Keep a Changelog](https://keepachangelog.com/en/1.0.0/) and
---
+## [Unreleased]
+
+### Added
+
+- **Web terminal capture** — WebSocket `/ws/session/{name}/terminal` now streams browser keystrokes to the PTY and logs each newline-delimited command as a `CommandEvent` (with cwd, seq, and part) into `session.jsonl`. Live terminal UI added to session pages with REST fallbacks for start/read/write/stop.
+
+---
+
## [0.5.0] — 2026-04-06
### Added
diff --git a/README.md b/README.md
index b56080e..e24e12c 100644
--- a/README.md
+++ b/README.md
@@ -372,6 +372,12 @@ Set `GUILD_SCROLL_ALLOW_REMOTE=1` to allow the report server to bind to non-loca
| `GET` | `/api/session/{name}/download` | Download session export (`?format=md\|html`) |
| `GET` | `/api/session/{name}/discoveries` | Fetch recent notes/assets timeline |
+### Live Web Terminal
+
+- Session pages include an **Open Terminal** panel that POSTs `/api/session/{name}/terminal/start` and streams over WebSocket `/ws/session/{name}/terminal`.
+- Each newline-delimited stdin message is logged to `session.jsonl` as a `CommandEvent` (with cwd, seq, and part) before being forwarded to the PTY; terminal output is mirrored to `terminal.log`.
+- REST fallbacks are available for polling and automation: `GET /api/session/{name}/terminal/read`, `POST /api/session/{name}/terminal/write`, and `POST /api/session/{name}/terminal/stop`.
+
### JSONL Event Types
| Type | Key Fields |
diff --git a/src/guild_scroll/web/app.py b/src/guild_scroll/web/app.py
index 0e904a5..032da70 100644
--- a/src/guild_scroll/web/app.py
+++ b/src/guild_scroll/web/app.py
@@ -1,7 +1,11 @@
from __future__ import annotations
+import base64
+import hashlib
import html
import json
+import os
+import queue
import re
import shutil
import socket
@@ -26,9 +30,17 @@
from guild_scroll.session_loader import LoadedSession, load_session
from guild_scroll.utils import generate_session_id, iso_timestamp, sanitize_session_name
from guild_scroll.validator import repair_session, validate_session
+from guild_scroll.web.terminal import (
+ TERMINALS,
+ ShellNotFound,
+ TerminalAlreadyRunning,
+ TerminalNotFound,
+ TerminalNotSupported,
+)
_SAFE_FILENAME_RE = re.compile(r"[^A-Za-z0-9._-]+")
+_WS_MAGIC = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"
def _write_jsonl_record(log_path: Path, record: dict[str, object]) -> None:
@@ -59,6 +71,19 @@ def _query_value(params: dict[str, list[str]], key: str) -> str | None:
return value or None
+def _parse_part(params: dict[str, list[str]]) -> int:
+ raw = _query_value(params, "part")
+ if raw is None:
+ return 1
+ try:
+ part = int(raw)
+ except ValueError as exc:
+ raise ValueError("part must be a positive integer") from exc
+ if part < 1:
+ raise ValueError("part must be a positive integer")
+ return part
+
+
def _parse_discovery_filters(params: dict[str, list[str]]) -> tuple[str | None, int]:
tag = _query_value(params, "tag")
limit_raw = _query_value(params, "limit")
@@ -562,11 +587,94 @@ def _render_session_page(
.discovery-summary {{ margin-top: 0.35rem; color: #edf4ff; word-break: break-word; }}
.discovery-tags {{ margin-top: 0.25rem; color: #9eb8da; font-size: 0.78rem; word-break: break-word; }}
.discovery-empty {{ margin: 0.8rem 0 0; color: #9eb8da; }}
+.terminal-panel {{ border: 1px solid #36567f; border-radius: 12px; background: linear-gradient(145deg, #101b2c, #0c1626); padding: 0.9rem; margin-bottom: 1rem; }}
+.terminal-header {{ display: flex; justify-content: space-between; align-items: center; gap: 0.5rem; }}
+.gs-terminal-btn {{ border: 1px solid #3d608d; background: #0e1a2c; color: #e9efff; padding: 0.4rem 0.8rem; border-radius: 8px; cursor: pointer; }}
+.gs-terminal-btn:hover {{ border-color: #52d0ff; background: #13243b; }}
+.terminal-output {{ background: #0b1020; border: 1px solid #334b70; border-radius: 8px; padding: 0.6rem; min-height: 220px; max-height: 320px; overflow: auto; white-space: pre-wrap; }}
+.terminal-actions {{ display: flex; gap: 0.5rem; align-items: center; margin-top: 0.55rem; }}
+.terminal-input {{ flex: 1; border: 1px solid #3d608d; background: #0e1a2c; color: #e9efff; padding: 0.45rem; border-radius: 6px; }}
@media (max-width: 980px) {{
.layout {{ grid-template-columns: 1fr; }}
.discoveries-panel {{ position: static; }}
}}
+
@@ -581,6 +689,18 @@ def _render_session_page(
+
+
{preview_markup}
@@ -617,6 +737,14 @@ def do_GET(self) -> None:
parsed = urlparse(self.path)
params = parse_qs(parsed.query)
+ if parsed.path.startswith("/ws/session/") and parsed.path.endswith("/terminal"):
+ session_name = parsed.path[len("/ws/session/"):-len("/terminal")].strip("/")
+ self._handle_terminal_websocket(session_name, params)
+ return
+ if parsed.path.startswith("/api/session/") and parsed.path.endswith("/terminal/read"):
+ session_name = parsed.path[len("/api/session/"):-len("/terminal/read")].strip("/")
+ self._handle_terminal_read(session_name, params)
+ return
if parsed.path == "/":
self._handle_index()
return
@@ -653,6 +781,18 @@ def do_POST(self) -> None:
parsed = urlparse(self.path)
params = parse_qs(parsed.query)
+ if parsed.path.startswith("/api/session/") and parsed.path.endswith("/terminal/start"):
+ session_name = parsed.path[len("/api/session/"):-len("/terminal/start")].strip("/")
+ self._handle_terminal_start(session_name, params)
+ return
+ if parsed.path.startswith("/api/session/") and parsed.path.endswith("/terminal/write"):
+ session_name = parsed.path[len("/api/session/"):-len("/terminal/write")].strip("/")
+ self._handle_terminal_write(session_name, params)
+ return
+ if parsed.path.startswith("/api/session/") and parsed.path.endswith("/terminal/stop"):
+ session_name = parsed.path[len("/api/session/"):-len("/terminal/stop")].strip("/")
+ self._handle_terminal_stop(session_name, params)
+ return
if parsed.path == "/api/sessions":
self._handle_create_session()
return
@@ -710,6 +850,233 @@ def _handle_session_api(self, raw_name: str, params: dict[str, list[str]]) -> No
}
)
+ def _handle_terminal_start(self, raw_name: str, params: dict[str, list[str]]) -> None:
+ session_name = unquote(raw_name)
+ if not _is_safe_session_name(session_name):
+ self._send_json({"error": "Invalid session name"}, status=400)
+ return
+ try:
+ part = _parse_part(params)
+ except ValueError as exc:
+ self._send_json({"error": str(exc)}, status=400)
+ return
+
+ try:
+ proc = TERMINALS.start(session_name, part=part)
+ except FileNotFoundError:
+ self._send_json({"error": "Session not found"}, status=404)
+ return
+ except TerminalAlreadyRunning:
+ self._send_json({"error": "Terminal already running"}, status=409)
+ return
+ except TerminalNotSupported as exc:
+ self._send_json({"error": str(exc)}, status=501)
+ return
+ except ShellNotFound as exc:
+ self._send_json({"error": str(exc)}, status=500)
+ return
+
+ self._send_json({"started": True, "pid": proc.pid, "part": part})
+
+ def _handle_terminal_read(self, raw_name: str, params: dict[str, list[str]]) -> None:
+ session_name = unquote(raw_name)
+ if not _is_safe_session_name(session_name):
+ self._send_json({"error": "Invalid session name"}, status=400)
+ return
+ try:
+ part = _parse_part(params)
+ except ValueError as exc:
+ self._send_json({"error": str(exc)}, status=400)
+ return
+
+ alive, output = TERMINALS.read(session_name, part=part)
+ self._send_json({"alive": bool(alive), "output": output})
+
+ def _handle_terminal_write(self, raw_name: str, params: dict[str, list[str]]) -> None:
+ session_name = unquote(raw_name)
+ if not _is_safe_session_name(session_name):
+ self._send_json({"error": "Invalid session name"}, status=400)
+ return
+ try:
+ part = _parse_part(params)
+ except ValueError as exc:
+ self._send_json({"error": str(exc)}, status=400)
+ return
+
+ body = self._read_json_body()
+ if not isinstance(body, dict):
+ self._send_json({"error": "Invalid request body"}, status=400)
+ return
+ payload = str(body.get("input", ""))
+ if not payload:
+ self._send_json({"error": "Input is required"}, status=400)
+ return
+
+ try:
+ TERMINALS.write(session_name, payload, part=part)
+ except TerminalNotFound:
+ self._send_json({"error": "No active terminal"}, status=404)
+ return
+ self._send_json({"ok": True})
+
+ def _handle_terminal_stop(self, raw_name: str, params: dict[str, list[str]]) -> None:
+ session_name = unquote(raw_name)
+ if not _is_safe_session_name(session_name):
+ self._send_json({"error": "Invalid session name"}, status=400)
+ return
+ try:
+ part = _parse_part(params)
+ except ValueError as exc:
+ self._send_json({"error": str(exc)}, status=400)
+ return
+
+ try:
+ TERMINALS.stop(session_name, part=part)
+ except TerminalNotFound:
+ self._send_json({"error": "No active terminal"}, status=404)
+ return
+ self._send_json({"stopped": True})
+
+ def _handle_terminal_websocket(self, raw_name: str, params: dict[str, list[str]]) -> None:
+ session_name = unquote(raw_name)
+ if not _is_safe_session_name(session_name):
+ self._send_text("Invalid session name", status=400)
+ return
+ try:
+ part = _parse_part(params)
+ except ValueError as exc:
+ self._send_text(str(exc), status=400)
+ return
+
+ terminal = TERMINALS.get(session_name, part=part)
+ if terminal is None:
+ self._send_text("No active terminal", status=404)
+ return
+
+ upgrade = self.headers.get("Upgrade", "").lower()
+ if upgrade != "websocket":
+ self._send_text("Upgrade header required", status=400)
+ return
+ key = self.headers.get("Sec-WebSocket-Key")
+ if not key:
+ self._send_text("Missing Sec-WebSocket-Key", status=400)
+ return
+
+ accept = base64.b64encode(
+ hashlib.sha1((key + _WS_MAGIC).encode("utf-8")).digest()
+ ).decode("utf-8")
+ self.send_response(101, "Switching Protocols")
+ self.send_header("Upgrade", "websocket")
+ self.send_header("Connection", "Upgrade")
+ self.send_header("Sec-WebSocket-Accept", accept)
+ self.end_headers()
+ self._serve_terminal_socket(terminal)
+
+ def _serve_terminal_socket(self, terminal) -> None:
+ conn = self.connection
+ conn.settimeout(0.2)
+ subscriber = terminal.add_subscriber()
+ try:
+ self._flush_terminal_output(subscriber)
+ while True:
+ try:
+ opcode, payload = self._read_ws_frame()
+ except TimeoutError:
+ if not terminal.is_alive():
+ break
+ self._flush_terminal_output(subscriber)
+ continue
+ if opcode is None:
+ break
+ if opcode == 0x8:
+ break
+ if opcode == 0x1:
+ text = payload.decode("utf-8", errors="replace")
+ try:
+ terminal.write(text)
+ except TerminalNotFound:
+ break
+ elif opcode == 0x9: # ping
+ self._send_ws_frame(b"", opcode=0xA)
+ self._flush_terminal_output(subscriber)
+ if not terminal.is_alive():
+ break
+ self._flush_terminal_output(subscriber)
+ finally:
+ terminal.remove_subscriber(subscriber)
+ try:
+ conn.close()
+ except Exception:
+ pass
+ self.close_connection = True
+
+ def _send_ws_frame(self, payload: bytes, opcode: int = 0x1) -> None:
+ length = len(payload)
+ frame = bytearray()
+ frame.append(0x80 | (opcode & 0x0F))
+ if length < 126:
+ frame.append(length)
+ elif length < (1 << 16):
+ frame.append(126)
+ frame.extend(length.to_bytes(2, "big"))
+ else:
+ frame.append(127)
+ frame.extend(length.to_bytes(8, "big"))
+ frame.extend(payload)
+ try:
+ self.connection.sendall(frame)
+ except Exception:
+ pass
+
+ def _read_ws_frame(self) -> tuple[int | None, bytes]:
+ try:
+ header = self.connection.recv(2)
+ except socket.timeout as exc:
+ raise TimeoutError from exc
+ if not header or len(header) < 2:
+ return None, b""
+ b1, b2 = header
+ opcode = b1 & 0x0F
+ masked = b2 & 0x80
+ length = b2 & 0x7F
+ if length == 126:
+ ext = self.connection.recv(2)
+ length = int.from_bytes(ext, "big")
+ elif length == 127:
+ ext = self.connection.recv(8)
+ length = int.from_bytes(ext, "big")
+
+ mask_key = b""
+ if masked:
+ mask_key = self.connection.recv(4)
+
+ payload = b""
+ remaining = length
+ while remaining > 0:
+ chunk = self.connection.recv(remaining)
+ if not chunk:
+ break
+ payload += chunk
+ remaining -= len(chunk)
+
+ if masked and mask_key:
+ payload = bytes(b ^ mask_key[i % 4] for i, b in enumerate(payload))
+ return opcode, payload
+
+ def _flush_terminal_output(self, subscriber: queue.SimpleQueue[str]) -> None:
+ while True:
+ try:
+ chunk = subscriber.get_nowait()
+ except queue.Empty:
+ break
+ if not chunk:
+ continue
+ if isinstance(chunk, str):
+ payload = chunk.encode("utf-8")
+ else:
+ payload = chunk
+ self._send_ws_frame(payload, opcode=0x1)
+
def _handle_report(self, raw_name: str, params: dict[str, list[str]]) -> None:
try:
session = self._load_filtered_session(raw_name, params)
diff --git a/src/guild_scroll/web/terminal.py b/src/guild_scroll/web/terminal.py
new file mode 100644
index 0000000..d6f58fe
--- /dev/null
+++ b/src/guild_scroll/web/terminal.py
@@ -0,0 +1,302 @@
+from __future__ import annotations
+
+import json
+import os
+import queue
+import shutil
+import subprocess
+import tempfile
+import threading
+import time
+from pathlib import Path
+from typing import Dict, Optional, Tuple
+
+from guild_scroll.config import PARTS_DIR_NAME, SESSION_LOG_NAME, get_sessions_dir
+from guild_scroll.integrity import load_session_key
+from guild_scroll.log_schema import CommandEvent
+from guild_scroll.log_writer import JSONLWriter
+from guild_scroll.session import _patch_session_meta, _read_session_meta, _session_root_from_logs_dir
+from guild_scroll.utils import iso_timestamp
+
+
+class TerminalError(Exception):
+ """Base error for terminal operations."""
+
+
+class TerminalNotSupported(TerminalError):
+ """Raised when PTY terminals are not supported on the current platform."""
+
+
+class TerminalAlreadyRunning(TerminalError):
+ """Raised when attempting to start a terminal that already exists."""
+
+
+class TerminalNotFound(TerminalError):
+ """Raised when operating on a non-existent terminal session."""
+
+
+class ShellNotFound(TerminalError):
+ """Raised when the configured shell is not available on the host."""
+
+
+def _session_paths(session_name: str, part: int) -> Tuple[Path, Path]:
+ """Return (session_dir, logs_dir) for the given session and part."""
+ sess_dir = get_sessions_dir() / session_name
+ if part <= 1:
+ logs_dir = sess_dir / "logs"
+ else:
+ logs_dir = sess_dir / PARTS_DIR_NAME / str(part) / "logs"
+ return sess_dir, logs_dir
+
+
+class TerminalProcess:
+ """Manage a single PTY-backed shell for a session."""
+
+ def __init__(self, session_name: str, part: int = 1, shell: str = "zsh"):
+ try:
+ import pty
+ import fcntl
+ except ImportError as exc: # pragma: no cover - platform specific
+ raise TerminalNotSupported("Terminal not supported on this platform") from exc
+
+ shell_path = shutil.which(shell)
+ if not shell_path:
+ raise ShellNotFound(f"{shell} not found on this system")
+
+ self.session_name = session_name
+ self.part = part
+ self._pty = pty
+ self._fcntl = fcntl
+ self._input_buffer = ""
+ self._alive = True
+ self._lock = threading.Lock()
+ self._output_buffer: list[str] = []
+ self._subscribers: set[queue.SimpleQueue[str]] = set()
+
+ sess_dir, logs_dir = _session_paths(session_name, part)
+ logs_dir.mkdir(parents=True, exist_ok=True)
+ sess_dir.mkdir(parents=True, exist_ok=True)
+ self._logs_dir = logs_dir
+ self._sess_dir = sess_dir
+ self._log_file = sess_dir / "terminal.log"
+ self._log_file.touch(exist_ok=True)
+
+ self._session_log_path = logs_dir / SESSION_LOG_NAME
+ self._session_root = _session_root_from_logs_dir(logs_dir, part)
+ self._hmac_key = load_session_key(self._session_root)
+ self._seq = self._next_seq()
+
+ master_fd, slave_fd = pty.openpty()
+ self._master_fd = master_fd
+ self._slave_fd = slave_fd
+
+ env = os.environ.copy()
+ env.setdefault("TERM", "xterm-256color")
+ work_dir = tempfile.gettempdir()
+ self._proc = subprocess.Popen(
+ [shell_path],
+ stdin=slave_fd,
+ stdout=slave_fd,
+ stderr=slave_fd,
+ cwd=work_dir,
+ env=env,
+ preexec_fn=os.setsid if hasattr(os, "setsid") else None,
+ )
+ self._fcntl.fcntl(master_fd, self._fcntl.F_SETFL, os.O_NONBLOCK)
+
+ self._reader = threading.Thread(target=self._read_loop, daemon=True)
+ self._reader.start()
+
+ @property
+ def pid(self) -> int:
+ return self._proc.pid
+
+ def is_alive(self) -> bool:
+ return self._alive and self._proc.poll() is None
+
+ def stop(self) -> None:
+ with self._lock:
+ if not self.is_alive():
+ self._alive = False
+ try:
+ if self._proc.poll() is None:
+ self._proc.terminate()
+ try:
+ self._proc.wait(timeout=1.5)
+ except subprocess.TimeoutExpired:
+ self._proc.kill()
+ finally:
+ try:
+ os.close(self._master_fd)
+ except OSError:
+ pass
+ try:
+ os.close(self._slave_fd)
+ except OSError:
+ pass
+ self._alive = False
+
+ def _read_loop(self) -> None:
+ while self.is_alive():
+ try:
+ data = os.read(self._master_fd, 4096)
+ except BlockingIOError:
+ time.sleep(0.05)
+ continue
+ except OSError:
+ break
+
+ if not data:
+ break
+
+ text = data.decode("utf-8", errors="replace")
+ with self._lock:
+ self._output_buffer.append(text)
+ try:
+ with self._log_file.open("a", encoding="utf-8") as fh:
+ fh.write(text)
+ except OSError:
+ pass
+
+ for subscriber in list(self._subscribers):
+ try:
+ subscriber.put_nowait(text)
+ except Exception:
+ self._subscribers.discard(subscriber)
+
+ self._alive = False
+
+ def _next_seq(self) -> int:
+ if not self._session_log_path.exists():
+ return 1
+ count = 0
+ try:
+ from guild_scroll.crypto import read_plaintext
+
+ for line in read_plaintext(self._session_log_path).splitlines():
+ line = line.strip()
+ if not line:
+ continue
+ try:
+ record = json.loads(line)
+ if record.get("type") == "command":
+ count += 1
+ except Exception:
+ continue
+ except Exception:
+ pass
+ return count + 1
+
+ def _capture_cwd(self) -> str:
+ try:
+ return os.readlink(f"/proc/{self._proc.pid}/cwd")
+ except OSError:
+ try:
+ result = subprocess.run(
+ ["pwd"], capture_output=True, text=True, check=False, timeout=1.0
+ )
+ return result.stdout.strip()
+ except Exception:
+ return ""
+
+ def _append_command_event(self, command: str) -> None:
+ ts = iso_timestamp()
+ event = CommandEvent(
+ seq=self._seq,
+ command=command,
+ timestamp_start=ts,
+ timestamp_end=ts,
+ exit_code=-1,
+ working_directory=self._capture_cwd(),
+ part=self.part,
+ )
+ writer = JSONLWriter(self._session_log_path, hmac_key=self._hmac_key)
+ writer.write(event.to_dict())
+ writer.close()
+
+ meta = _read_session_meta(self._session_log_path)
+ end_time = meta.get("end_time") if meta else None
+ _patch_session_meta(self._session_log_path, end_time, self._seq)
+ self._seq += 1
+
+ def write(self, data: str) -> None:
+ if not self.is_alive():
+ raise TerminalNotFound("Terminal is not running")
+
+ self._input_buffer += data
+ lines = self._input_buffer.split("\n")
+ self._input_buffer = lines.pop() # Remaining partial line
+ for line in lines:
+ candidate = line.rstrip("\r")
+ if candidate.strip():
+ self._append_command_event(candidate)
+
+ os.write(self._master_fd, data.encode("utf-8"))
+
+ def read_output(self) -> Tuple[bool, str]:
+ with self._lock:
+ combined = "".join(self._output_buffer)
+ self._output_buffer.clear()
+ return self.is_alive(), combined
+
+ def add_subscriber(self) -> queue.SimpleQueue[str]:
+ subscriber: queue.SimpleQueue[str] = queue.SimpleQueue()
+ self._subscribers.add(subscriber)
+ return subscriber
+
+ def remove_subscriber(self, subscriber: queue.SimpleQueue[str]) -> None:
+ self._subscribers.discard(subscriber)
+
+
+class TerminalManager:
+ """Thread-safe registry of PTY terminals keyed by session and part."""
+
+ def __init__(self) -> None:
+ self._sessions: Dict[Tuple[str, int], TerminalProcess] = {}
+ self._lock = threading.Lock()
+
+ def start(self, session_name: str, part: int = 1) -> TerminalProcess:
+ sess_dir, logs_dir = _session_paths(session_name, part)
+ if not logs_dir.exists():
+ raise FileNotFoundError(f"Session not found: {session_name!r}")
+
+ key = (session_name, part)
+ with self._lock:
+ existing = self._sessions.get(key)
+ if existing and existing.is_alive():
+ raise TerminalAlreadyRunning(f"Terminal already active for {session_name!r}")
+ process = TerminalProcess(session_name=session_name, part=part)
+ self._sessions[key] = process
+ return process
+
+ def get(self, session_name: str, part: int = 1) -> Optional[TerminalProcess]:
+ with self._lock:
+ proc = self._sessions.get((session_name, part))
+ if proc and proc.is_alive():
+ return proc
+ return None
+
+ def stop(self, session_name: str, part: int = 1) -> None:
+ key = (session_name, part)
+ with self._lock:
+ proc = self._sessions.get(key)
+ if not proc:
+ raise TerminalNotFound(f"No active terminal for {session_name!r}")
+ proc.stop()
+ with self._lock:
+ self._sessions.pop(key, None)
+
+ def read(self, session_name: str, part: int = 1) -> Tuple[bool, str]:
+ proc = self.get(session_name, part)
+ if not proc:
+ return False, ""
+ return proc.read_output()
+
+ def write(self, session_name: str, data: str, part: int = 1) -> None:
+ proc = self.get(session_name, part)
+ if not proc:
+ raise TerminalNotFound(f"No active terminal for {session_name!r}")
+ proc.write(data)
+
+
+TERMINALS = TerminalManager()
diff --git a/tests/test_web.py b/tests/test_web.py
index 840342c..1c2623a 100644
--- a/tests/test_web.py
+++ b/tests/test_web.py
@@ -1,7 +1,10 @@
+import base64
import errno
import json
+import os
import random
import re
+import socket
import ssl
import string
import sys
@@ -130,6 +133,71 @@ def _multipart_body(filename: str, data: bytes, field: str = "file") -> tuple[by
return body, f"multipart/form-data; boundary={boundary}"
+def _ws_connect(server, path: str) -> socket.socket:
+ host = "127.0.0.1"
+ port = server.server_port
+ sock = socket.create_connection((host, port))
+ key = base64.b64encode(os.urandom(16)).decode()
+ headers = (
+ f"GET {path} HTTP/1.1\r\n"
+ f"Host: {host}:{port}\r\n"
+ "Upgrade: websocket\r\n"
+ "Connection: Upgrade\r\n"
+ f"Sec-WebSocket-Key: {key}\r\n"
+ "Sec-WebSocket-Version: 13\r\n"
+ "\r\n"
+ )
+ sock.sendall(headers.encode())
+ response = sock.recv(1024)
+ assert b"101" in response
+ return sock
+
+
+def _ws_send_text(sock: socket.socket, message: str) -> None:
+ data = message.encode("utf-8")
+ length = len(data)
+ frame = bytearray()
+ frame.append(0x81) # FIN + text frame
+ mask_key = os.urandom(4)
+ if length < 126:
+ frame.append(0x80 | length)
+ elif length < (1 << 16):
+ frame.append(0x80 | 126)
+ frame.extend(length.to_bytes(2, "big"))
+ else:
+ frame.append(0x80 | 127)
+ frame.extend(length.to_bytes(8, "big"))
+ masked = bytes(b ^ mask_key[i % 4] for i, b in enumerate(data))
+ frame.extend(mask_key)
+ frame.extend(masked)
+ sock.sendall(frame)
+
+
+def _ws_recv_text(sock: socket.socket, timeout: float = 1.5) -> str:
+ sock.settimeout(timeout)
+ header = sock.recv(2)
+ if not header or len(header) < 2:
+ raise socket.timeout()
+ b1, b2 = header
+ opcode = b1 & 0x0F
+ if opcode == 0x8: # close
+ return ""
+ length = b2 & 0x7F
+ if length == 126:
+ length = int.from_bytes(sock.recv(2), "big")
+ elif length == 127:
+ length = int.from_bytes(sock.recv(8), "big")
+ payload = b""
+ remaining = length
+ while remaining > 0:
+ chunk = sock.recv(remaining)
+ if not chunk:
+ break
+ payload += chunk
+ remaining -= len(chunk)
+ return payload.decode("utf-8", errors="replace")
+
+
class TestIsSafeSessionName:
def test_rejects_traversal_strings(self, isolated_sessions_dir):
assert not _is_safe_session_name("../escape")
@@ -1426,6 +1494,65 @@ def test_terminal_full_lifecycle(self, isolated_sessions_dir):
log_path = sessions_dir / "term-full" / "terminal.log"
assert log_path.exists()
+ @pytest.mark.skipif(sys.platform == "win32", reason="pty not available on Windows")
+ def test_terminal_websocket_captures_input_as_command_event(self, isolated_sessions_dir):
+ session_name = "ws-capture"
+ sessions_dir = get_sessions_dir()
+ sessions_dir.mkdir(parents=True, exist_ok=True)
+ _make_session(sessions_dir, session_name)
+
+ with _running_server() as server:
+ start_status, _, start_resp = _request_post(
+ server, f"/api/session/{session_name}/terminal/start", b""
+ )
+ start_payload = json.loads(start_resp)
+ if start_status == 501:
+ pytest.skip("pty not available")
+ if start_status == 500 and "zsh not found" in start_payload.get("error", ""):
+ pytest.skip("zsh not available")
+ assert start_status == 200
+
+ ws_path = f"/ws/session/{quote(session_name)}/terminal"
+ sock = _ws_connect(server, ws_path)
+ try:
+ _ws_send_text(sock, "echo websocket-cmd\n")
+ output_seen = ""
+ for _ in range(20):
+ try:
+ output_seen += _ws_recv_text(sock, timeout=0.3)
+ except socket.timeout:
+ pass
+ if "websocket-cmd" in output_seen:
+ break
+ time.sleep(0.1)
+ finally:
+ try:
+ sock.close()
+ except Exception:
+ pass
+
+ _request_post(server, f"/api/session/{session_name}/terminal/stop", b"")
+
+ log_path = sessions_dir / session_name / "logs" / SESSION_LOG_NAME
+ records = [
+ json.loads(line)
+ for line in log_path.read_text(encoding="utf-8").splitlines()
+ if line.strip()
+ ]
+ commands = [
+ r for r in records if r.get("type") == "command" and r.get("command") == "echo websocket-cmd"
+ ]
+ assert commands
+ event = commands[-1]
+ assert event.get("part", 1) == 1
+ assert event.get("working_directory", "") != ""
+
+ with _running_server() as server:
+ status, _, body = _request(server, f"/api/session/{session_name}")
+ assert status == 200
+ api_payload = json.loads(body)
+ assert any(cmd.get("command") == "echo websocket-cmd" for cmd in api_payload["commands"])
+
# ── GET /api/sessions ─────────────────────────────────────────────────────────