Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
258 changes: 254 additions & 4 deletions hermes_cli/logs.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,15 @@
hermes logs --since 30m -f # follow, starting 30 min ago
"""

import io
import json
import os
import re
import stat
import sys
import time
from datetime import datetime, timedelta
from collections import deque
from pathlib import Path
from typing import Optional, Sequence

Expand All @@ -35,8 +40,7 @@
"gateway": "gateway.log",
"gui": "gui.log",
"desktop": "desktop.log",
# Every stdio MCP subprocess's stderr (tools/mcp_tool.py redirects it
# here, with per-server session markers) — the "MCP output channel".
# Legacy raw stderr and the index of profile-owned stdio attempt files.
"mcp": "mcp-stderr.log",
}

Expand Down Expand Up @@ -215,7 +219,8 @@ def tail_log(

# Read and display the tail
try:
lines = _read_tail(log_path, num_lines, has_filters=has_filters,
read_tail = _read_mcp_tail if log_name == "mcp" else _read_tail
lines = read_tail(log_path, num_lines, has_filters=has_filters,
min_level=min_level, session_filter=session,
since=since_dt, component_prefixes=component_prefixes)
except PermissionError:
Expand Down Expand Up @@ -247,12 +252,257 @@ def tail_log(

# Follow mode — poll for new content
try:
_follow_log(log_path, min_level=min_level, session_filter=session,
follow_log = _follow_mcp_log if log_name == "mcp" else _follow_log
follow_log(log_path, min_level=min_level, session_filter=session,
since=since_dt, component_prefixes=component_prefixes)
except KeyboardInterrupt:
print("\n--- stopped ---")


_MCP_MAX_ATTEMPTS = 128
_MCP_INDEX_LINES = 2000
_MCP_TAIL_BYTES = 1048576
_MCP_MAX_PREVIEW_LINES = 10000
_MCP_FOLLOW_BYTES = 65536
_MCP_MAX_LINE_BYTES = 65536


class _MCPPreviewNotice(str):
"""Reader metadata, distinct from child content subject to log filters."""


def _mcp_open_regular(path: Path):
return open(path, "rb", opener=lambda name, flags: os.open(
name, flags | getattr(os, "O_NOFOLLOW", 0) | getattr(os, "O_NONBLOCK", 0)
))


def _mcp_oversized_line(index: bool) -> str:
source = "index" if index else "stderr"
return _MCPPreviewNotice(f"MCP {source} line exceeds MCP preview limit; full content remains in the log file.\n")


def _mcp_tail_rows(path: Path, n: int, *, index: bool = False) -> list:
"""A finite byte window; the generic log reader is unchanged."""
if n <= 0:
return []
rows = deque(maxlen=min(n, _MCP_MAX_PREVIEW_LINES))
try:
with _mcp_open_regular(path) as stream:
info = os.fstat(stream.fileno())
if not stat.S_ISREG(info.st_mode):
return []
start = max(0, info.st_size - _MCP_TAIL_BYTES)
stream.seek(start)
data = stream.read(_MCP_TAIL_BYTES)
if start:
fragment, separator, data = data.partition(b"\n")
if not separator or len(fragment) > _MCP_MAX_LINE_BYTES:
rows.append(_mcp_oversized_line(index))
for row in io.BytesIO(data):
payload = row.removesuffix(b"\n").removesuffix(b"\r")
if len(payload) > _MCP_MAX_LINE_BYTES:
rows.append(_mcp_oversized_line(index))
else:
decoded = row.decode("utf-8", errors="replace")
# An EOF preview is a display record, even before the child
# finishes its line. The stored bytes remain unchanged.
rows.append(decoded if decoded.endswith("\n") else decoded + "\n")
except OSError:
return []
return list(rows)


def _mcp_control_record(line: str) -> Optional[dict]:
try:
record = json.loads(line)
except (ValueError, TypeError):
return None
return record if isinstance(record, dict) and record.get("kind") == "mcp.stdio.attempt" else None


def _mcp_attempt(record: dict, index: Path) -> Optional[dict]:
"""Resolve only this profile's UUID-named capture, never an arbitrary path."""
attempt = record.get("attempt_id")
if (not isinstance(attempt, str) or re.fullmatch(r"[0-9a-f]{32}", attempt) is None
or not isinstance(record.get("server"), str)
or type(record.get("parent_pid")) is not int or record["parent_pid"] <= 0
or record.get("destination") != "file"
or not isinstance(record.get("config_home"), str)
or not isinstance(record.get("stderr_path"), str)):
return None
try:
home = index.parent.parent.resolve()
expected = home / "logs" / "mcp-stderr" / f"{attempt}.log"
if (Path(record["config_home"]).resolve() != home
or Path(record["stderr_path"]).resolve() != expected
or not expected.is_file()):
return None
except (OSError, ValueError, RuntimeError):
return None
return {"attempt_id": attempt, "server": record["server"], "parent_pid": record["parent_pid"],
"config_home": str(home), "stderr_path": str(expected), "destination": "file"}


def _mcp_child_line(line: str, record: dict, filters: dict) -> Optional[str]:
control = _mcp_control_record(line)
if control is not None and control.get("attempt_id") == record["attempt_id"]:
return None
owner = " ".join(f"{key}={json.dumps(value)}" for key, value in (
("profile", record["config_home"]), ("server", record["server"]),
("attempt", record["attempt_id"]),
))
rendered = f"[mcp {owner}] {line}"
session = filters.get("session_filter")
if session is not None and session not in rendered:
return None
if isinstance(line, _MCPPreviewNotice):
return rendered
if not _matches_filters(line, **{**filters, "session_filter": None}):
return None
return rendered


def _read_mcp_tail(path: Path, num_lines: int, *, has_filters: bool = False, **filters) -> list:
"""Read actual stderr in attempt order; legacy rows keep their own format."""
if num_lines <= 0:
return []
requested = num_lines
num_lines = min(num_lines, _MCP_MAX_PREVIEW_LINES)
result = []
remaining = _MCP_TAIL_BYTES
exhausted = False
seen = set()

def append(line):
nonlocal remaining, exhausted
if len(line) > remaining:
exhausted = True
return
result.append(line)
remaining -= len(line)

for line in reversed(_mcp_tail_rows(path, max(num_lines * 20, _MCP_INDEX_LINES), index=True)):
control = _mcp_control_record(line)
if control is None:
if isinstance(line, _MCPPreviewNotice) or _matches_filters(line, **filters):
append(line)
else:
record = _mcp_attempt(control, path)
if record is None or record["attempt_id"] in seen or len(seen) >= _MCP_MAX_ATTEMPTS:
continue
seen.add(record["attempt_id"])
rows = _mcp_tail_rows(Path(record["stderr_path"]),
max(num_lines * 20, _MCP_INDEX_LINES) if has_filters else num_lines + 4)
for row in reversed(rows):
rendered = _mcp_child_line(row, record, filters)
if rendered is not None:
append(rendered)
if exhausted or len(result) >= num_lines:
break
if exhausted or len(result) >= num_lines:
break
lines = list(reversed(result))
if exhausted or requested > _MCP_MAX_PREVIEW_LINES:
lines.insert(0, "MCP preview output limit reached; full content remains in the log files.\n")
return lines


def _mcp_cursor(path: Path):
try:
with _mcp_open_regular(path) as stream:
info = os.fstat(stream.fileno())
if not stat.S_ISREG(info.st_mode):
return None
# Seed only the bounded final fragment. If it later completes,
# display its complete new record; never replay complete history.
start = max(0, info.st_size - _MCP_MAX_LINE_BYTES - 1)
stream.seek(start)
recent = stream.read(min(info.st_size, _MCP_MAX_LINE_BYTES + 1))
pending = recent.rsplit(b"\n", 1)[-1]
dropping = len(pending) > _MCP_MAX_LINE_BYTES
return ((info.st_dev, info.st_ino), info.st_size, b"" if dropping else pending, dropping)
except OSError:
return None


def _mcp_read_chunk(path: Path, cursor, *, index: bool = False):
"""Frame bounded binary chunks; retain fragments, never file descriptors."""
try:
with _mcp_open_regular(path) as stream:
info = os.fstat(stream.fileno())
if not stat.S_ISREG(info.st_mode):
return [], cursor
identity = (info.st_dev, info.st_ino)
retained = cursor is not None and cursor[0] == identity and cursor[1] <= info.st_size
offset, pending, dropping = cursor[1:] if retained else (0, b"", False)
stream.seek(offset)
chunk = stream.read(_MCP_FOLLOW_BYTES)
offset = stream.tell()
data = pending + chunk
lines = []
if dropping:
_fragment, separator, data = data.partition(b"\n")
if not separator:
return [], (identity, offset, b"", True)
parts = data.split(b"\n")
for row in parts[:-1]:
lines.append(_mcp_oversized_line(index) if len(row) > _MCP_MAX_LINE_BYTES
else row.decode("utf-8", errors="replace") + "\n")
pending = parts[-1]
dropping = len(pending) > _MCP_MAX_LINE_BYTES
if dropping:
lines.append(_mcp_oversized_line(index))
pending = b""
return lines, (identity, offset, pending, dropping)
except OSError:
return [], cursor


def _follow_mcp_log(path: Path, **filters) -> None:
records = {}

def remember(line):
control = _mcp_control_record(line)
record = _mcp_attempt(control, path) if control is not None else None
if record is not None:
records[record["attempt_id"]] = record
if len(records) > _MCP_MAX_ATTEMPTS:
records.pop(next(iter(records)))
print("MCP follow preview retains the latest 128 attempts; older captures remain in the log files.")

for line in _mcp_tail_rows(path, _MCP_INDEX_LINES, index=True):
remember(line)
positions = {key: _mcp_cursor(Path(record["stderr_path"])) for key, record in records.items()}
for key, cursor in positions.items():
if cursor is not None and cursor[3]:
rendered = _mcp_child_line(_mcp_oversized_line(False), records[key], filters)
if rendered is not None:
print(rendered, end="")
index_cursor = _mcp_cursor(path)
if index_cursor is not None and index_cursor[3]:
print(_mcp_oversized_line(True), end="")
while True:
lines, index_cursor = _mcp_read_chunk(path, index_cursor, index=True)
for line in lines:
if _mcp_control_record(line) is not None:
remember(line)
elif isinstance(line, _MCPPreviewNotice) or _matches_filters(line, **filters):
print(line, end="")
sys.stdout.flush()
for key, record in records.items():
if _mcp_attempt(record, path) is None:
continue
lines, positions[key] = _mcp_read_chunk(Path(record["stderr_path"]), positions.get(key))
for line in lines:
rendered = _mcp_child_line(line, record, filters)
if rendered is not None:
print(rendered, end="")
sys.stdout.flush()
positions = {key: positions[key] for key in records if key in positions}
time.sleep(0.3)


def _read_tail(
path: Path,
num_lines: int,
Expand Down
Loading
Loading