refactor(gateway): stream_consumer*/slash_commands* — hug closing brackets (AST-identical)

This commit is contained in:
Teknium
2026-09-02 23:42:01 -07:00
parent d6268539d4
commit e530f5936f
5 changed files with 102 additions and 205 deletions
+32 -64
View File
@@ -28,8 +28,7 @@ from gateway.session import AsyncSessionStore
from gateway.slash_commands_goals import GatewayGoalCommandsMixin
from gateway.slash_commands_model import ( # noqa: F401 — _model_switch_skew_guard re-exported for tests
GatewayModelCommandsMixin,
_model_switch_skew_guard,
)
_model_switch_skew_guard)
from gateway.slash_commands_session import GatewaySessionCommandsMixin
from gateway.slash_commands_status import GatewayStatusCommandsMixin
from hermes_cli.config import atomic_config_write, cfg_get
@@ -42,8 +41,7 @@ logger = logging.getLogger("gateway.run")
_ROLLBACK_SKIP_LINES = (
("skipped_user_edits", "gateway.rollback.kept_user_edits"),
("skipped_oversize", "gateway.rollback.kept_oversize"),
("failed_deletes", "gateway.rollback.failed_deletes"),
)
("failed_deletes", "gateway.rollback.failed_deletes"))
# /busy input modes -> (status-card behavior, set-confirmation behavior).
_BUSY_MODE_BEHAVIOR = {
@@ -57,34 +55,29 @@ _BUSY_MODE_BEHAVIOR = {
_DIFF_MODE_BY_ARG = {
"staged": "staged", "--staged": "staged", "cached": "staged", "--cached": "staged",
"all": "all", "--all": "all", "head": "all",
"session": "session",
}
"session": "session"}
# /voice subcommand -> stored mode (None = auto-TTS disabled), confirmation i18n key.
_VOICE_MODE_BY_ARG = {
**dict.fromkeys(("on", "enable"), ("voice_only", "gateway.voice.enabled_voice_only")),
**dict.fromkeys(("off", "disable"), ("off", "gateway.voice.disabled_text")),
"tts": ("all", "gateway.voice.tts_enabled"),
}
"tts": ("all", "gateway.voice.tts_enabled")}
# /footer argument -> new enabled state ("" toggles; anything else is a usage error).
_FOOTER_STATE_BY_ARG = {
**dict.fromkeys(("on", "enable", "true", "1"), True),
**dict.fromkeys(("off", "disable", "false", "0"), False),
}
**dict.fromkeys(("off", "disable", "false", "0"), False)}
# /approve modifier tokens -> approval choice (default "once").
_APPROVE_CHOICE_BY_ARG = {
**dict.fromkeys(("always", "permanent", "permanently"), "always"),
**dict.fromkeys(("session", "ses"), "session"),
}
**dict.fromkeys(("session", "ses"), "session")}
_PLATFORM_USAGE = (
"Usage: /platform <list|pause|resume> [name]\n"
" /platform list — show platform status\n"
" /platform pause <name> — stop retrying a failing platform\n"
" /platform resume <name> — re-queue a paused platform"
)
" /platform resume <name> — re-queue a paused platform")
_WINDOWS_UPDATE_HELPER = """
import os, subprocess, sys
@@ -155,8 +148,7 @@ def _spawn_detached_update(hermes_cmd, output_path, exit_code_path) -> None:
subprocess.Popen(
[sys.executable, "-c", _WINDOWS_UPDATE_HELPER, str(output_path), str(exit_code_path),
sys.executable, "-m", "hermes_cli.main", "update", "--gateway"],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, **windows_detach_popen_kwargs(),
)
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, **windows_detach_popen_kwargs())
return
hermes_cmd_str = " ".join(shlex.quote(part) for part in hermes_cmd)
update_cmd = (
@@ -164,8 +156,7 @@ def _spawn_detached_update(hermes_cmd, output_path, exit_code_path) -> None:
f" > {shlex.quote(str(output_path))} 2>&1; "
# Avoid `status=$?`: `status` is read-only in zsh and this template is reused in
# macOS/zsh operator wrappers, so keep it zsh-safe even though bash runs it here.
f"rc=$?; printf '%s' \"$rc\" > {shlex.quote(str(exit_code_path))}"
)
f"rc=$?; printf '%s' \"$rc\" > {shlex.quote(str(exit_code_path))}")
# Preferred: setsid creates a new session, fully detached; fallback start_new_session=True
# calls os.setsid() in the child.
setsid_bin = shutil.which("setsid")
@@ -192,8 +183,7 @@ class GatewaySlashCommandsMixin(
GatewayModelCommandsMixin,
GatewaySessionCommandsMixin,
GatewayStatusCommandsMixin,
GatewayGoalCommandsMixin,
):
GatewayGoalCommandsMixin):
"""In-session slash-command handlers for GatewayRunner (plus the helpers the sibling mixins share)."""
async_session_store: AsyncSessionStore
@@ -282,8 +272,7 @@ class GatewaySlashCommandsMixin(
try:
await adapter.send(
source.chat_id, confirmation_text, reply_to=event.message_id,
metadata={"is_approval_prompt": True, "force_proactive_send": True},
)
metadata={"is_approval_prompt": True, "force_proactive_send": True})
except Exception as exc:
logger.warning(
"Failed to send /%s confirmation to %s: %s", verb, source.chat_id, exc, exc_info=True,
@@ -331,8 +320,7 @@ class GatewaySlashCommandsMixin(
reply = _execute("profile", options={"profile_name": profile_name, "home_display": display})
return "\n".join([
t("gateway.profile.header", profile=reply.data["profile"]),
t("gateway.profile.home", home=reply.data["home"]),
])
t("gateway.profile.home", home=reply.data["home"])])
async def _handle_whoami_command(self, event: MessageEvent) -> str:
"""Handle /whoami — platform, DM-vs-group scope, tier and runnable commands (always allowed)."""
@@ -424,8 +412,7 @@ class GatewaySlashCommandsMixin(
user_id_alt=_field("user_id_alt"),
notifier_profile=getattr(self, "_kanban_notifier_profile", None) or self._active_profile_name(),
# Subscribing from chat: deliver the passive message and wake the destination agent.
delivery_mode="notify+wake", delivery_metadata=delivery_metadata,
)
delivery_mode="notify+wake", delivery_metadata=delivery_metadata)
finally:
conn.close()
await asyncio.to_thread(_sub)
@@ -465,8 +452,7 @@ class GatewaySlashCommandsMixin(
await _stop(sibling_key, "stop_command_thread_sibling")
logger.info(
"STOP (thread sibling) by %s — interrupted %d run(s) in thread: %s",
session_key, len(sibling_keys), ", ".join(sibling_keys),
)
session_key, len(sibling_keys), ", ".join(sibling_keys))
return EphemeralReply(t("gateway.stop.stopped"))
# No running agent anywhere for this scope. A platform status indicator can still be stuck —
@@ -536,8 +522,7 @@ class GatewaySlashCommandsMixin(
"Ignoring redelivered /restart (platform=%s, update_id=%s) — "
"already processed by a previous gateway instance.",
event.source.platform.value if event.source and event.source.platform else "?",
event.platform_update_id,
)
event.platform_update_id)
return ""
if self._restart_requested or self._draining:
count = self._running_agent_count()
@@ -607,16 +592,14 @@ class GatewaySlashCommandsMixin(
source.platform in {None, Platform.LOCAL, Platform.RELAY}
or not getattr(source, "user_id", None)
or not callable(fronts_platform)
or not fronts_platform(source.platform)
):
or not fronts_platform(source.platform)):
return t("gateway.set_home.save_failed",
error="Relay does not authenticate this logical home target")
thread_id = _home_thread_from_source(source)
home = HomeChannel(
platform=source.platform, chat_id=str(chat_id), name=chat_name, thread_id=thread_id,
user_id=str(source.user_id) if getattr(source, "user_id", None) else None,
scope_id=str(source.scope_id) if getattr(source, "scope_id", None) else None,
)
scope_id=str(source.scope_id) if getattr(source, "scope_id", None) else None)
# config.yaml is canonical because it can persist the authenticated logical-target
# provenance required by Relay after a restart.
try:
@@ -797,8 +780,7 @@ class GatewaySlashCommandsMixin(
prompt, event.source, task_id, event_message_id=self._reply_anchor_for_event(event),
# Forward image/audio attachments so the background agent can see them.
media_urls=list(event.media_urls) if event.media_urls else [],
media_types=list(event.media_types) if event.media_types else [],
))
media_types=list(event.media_types) if event.media_types else []))
return t("gateway.background.started", preview=_preview(prompt), task_id=task_id)
async def _handle_btw_command(self, event: MessageEvent) -> str:
@@ -837,8 +819,7 @@ class GatewaySlashCommandsMixin(
try:
answer = await asyncio.to_thread(
answer_side_question, question, history_snapshot,
parent_agent=parent_agent, main_runtime=main_runtime,
)
parent_agent=parent_agent, main_runtime=main_runtime)
reply = t("gateway.btw.answer", preview=preview, answer=answer or "")
except Exception as e:
logger.warning("/btw side question failed: %s", e)
@@ -859,8 +840,7 @@ class GatewaySlashCommandsMixin(
# the store persists to the same MEMORY/USER.md and honors the configured char limits).
out = handle_pending_subcommand(
wa.MEMORY, event.get_command_args().strip().split(), memory_store=load_on_disk_store(),
set_mode_fn=self._write_approval_setter("memory", event),
)
set_mode_fn=self._write_approval_setter("memory", event))
return out if out is not None else (
"Unknown /memory subcommand. Use: pending, approve <id>, reject <id>, approval <on|off>."
)
@@ -879,8 +859,7 @@ class GatewaySlashCommandsMixin(
"Enable it with /skills approval on, then review staged "
"writes here with /skills pending.")
out = handle_pending_subcommand(
wa.SKILLS, args, set_mode_fn=self._write_approval_setter("skills", event)
)
wa.SKILLS, args, set_mode_fn=self._write_approval_setter("skills", event))
if out is None:
return ("Unknown /skills subcommand on this platform. Use: pending, "
"approve <id>, reject <id>, diff <id>, approval <on|off>. "
@@ -929,8 +908,7 @@ class GatewaySlashCommandsMixin(
try:
user_config = _load_gateway_config()
gate_enabled = is_truthy_value(
cfg_get(user_config, "display", "tool_progress_command"), default=False
)
cfg_get(user_config, "display", "tool_progress_command"), default=False)
except Exception:
gate_enabled = False
if not gate_enabled:
@@ -958,13 +936,11 @@ class GatewaySlashCommandsMixin(
return EphemeralReply(
f"**Busy input mode: `{mode}`" + "\n"
f"Messages while busy: _{behavior}_" + "\n"
f"Change with `/busy queue`, `/busy steer`, or `/busy interrupt`."
)
f"Change with `/busy queue`, `/busy steer`, or `/busy interrupt`.")
if arg not in _BUSY_MODE_BEHAVIOR:
return EphemeralReply(
f"Unknown mode `{arg}`. Use `/busy queue`, `/busy steer`, or `/busy interrupt`."
)
f"Unknown mode `{arg}`. Use `/busy queue`, `/busy steer`, or `/busy interrupt`.")
# Persist before mutate
from cli import save_config_value
@@ -1029,8 +1005,7 @@ class GatewaySlashCommandsMixin(
# Show a preview using current agent state if available.
preview = format_runtime_footer(
model=_resolve_gateway_model(user_config) or None, context_tokens=0, context_length=None,
fields=effective.get("fields") or ["model", "context_pct", "cwd"],
)
fields=effective.get("fields") or ["model", "context_pct", "cwd"])
if preview:
example = t("gateway.footer.example_line", preview=preview)
return t("gateway.footer.saved", state=_state(new_state), example=example)
@@ -1067,8 +1042,7 @@ class GatewaySlashCommandsMixin(
return result
return await self._request_slash_confirm(
event=event, command="reload-mcp", title="/reload-mcp",
message=t("gateway.reload_mcp.confirm_prompt"), handler=_on_confirm,
)
message=t("gateway.reload_mcp.confirm_prompt"), handler=_on_confirm)
async def _handle_reload_skills_command(self, event: MessageEvent) -> str:
"""Handle /reload-skills — rescan skills dir, queue a note for next turn. Skills are invoked at
@@ -1114,8 +1088,7 @@ class GatewaySlashCommandsMixin(
sections = ["[USER INITIATED SKILLS RELOAD:"]
for i18n_key, note_header, items in (
("gateway.reload_skills.added_header", "Added Skills:", added),
("gateway.reload_skills.removed_header", "Removed Skills:", removed),
):
("gateway.reload_skills.removed_header", "Removed Skills:", removed)):
if items:
formatted = [_fmt_line(item) for item in items]
lines += [t(i18n_key)] + formatted
@@ -1145,8 +1118,7 @@ class GatewaySlashCommandsMixin(
"No skill bundles installed.\n"
"Create one on the host with:\n"
" `hermes bundles create <name> --skill <s1> --skill <s2>`\n"
f"Directory: `{reply.data['dir']}`"
)
f"Directory: `{reply.data['dir']}`")
lines = [f"**Skill Bundles** ({len(bundles)} installed):", ""]
for info in bundles:
@@ -1172,8 +1144,7 @@ class GatewaySlashCommandsMixin(
signalling the event resumes them so the command executes inline (same flow as the CLI)."""
from tools.approval import resolve_gateway_approval
session_key, stale = self._blocking_approval_or_stale(
event, "gateway.approval_expired", "gateway.approve.no_pending"
)
event, "gateway.approval_expired", "gateway.approve.no_pending")
if stale:
return stale
@@ -1193,8 +1164,7 @@ class GatewaySlashCommandsMixin(
the CLI. ``/deny`` denies the oldest; ``/deny all`` denies everything."""
from tools.approval import resolve_gateway_approval
session_key, stale = self._blocking_approval_or_stale(
event, "gateway.deny.stale", "gateway.deny.no_pending"
)
event, "gateway.deny.stale", "gateway.deny.no_pending")
if stale:
return stale
@@ -1219,8 +1189,7 @@ class GatewaySlashCommandsMixin(
protect privacy; ``hermes debug share`` from the CLI does full uploads."""
from hermes_cli.debug import (
_GATEWAY_PRIVACY_NOTICE, _best_effort_sweep_expired_pastes, _capture_dump, _schedule_auto_delete,
collect_debug_report, upload_to_pastebin,
)
collect_debug_report, upload_to_pastebin)
# Run blocking I/O (dump capture, log reads, uploads) in a thread.
def _collect_and_upload():
@@ -1274,8 +1243,7 @@ class GatewaySlashCommandsMixin(
pending = {
"platform": src.platform.value, "chat_id": src.chat_id, "chat_type": src.chat_type,
"user_id": src.user_id, "session_key": self._session_key_for_source(src),
"timestamp": datetime.now().isoformat(),
}
"timestamp": datetime.now().isoformat()}
pending.update({k: v for k, v in (("thread_id", src.thread_id), ("message_id", event.message_id)) if v})
_tmp_pending = pending_path.with_suffix(".tmp")
_tmp_pending.write_text(json.dumps(pending), encoding="utf-8")
+51 -103
View File
@@ -32,8 +32,7 @@ _DM_CHAT_TYPES = {"dm", "direct", "private", ""}
_BRANCH_COPIED_FIELDS = (
"content", "tool_calls", "tool_call_id", "finish_reason", "reasoning", "reasoning_content",
"reasoning_details", "codex_reasoning_items", "codex_message_items", "timestamp",
)
"reasoning_details", "codex_reasoning_items", "codex_message_items", "timestamp")
def _sattr(obj, name: str) -> str:
@@ -65,8 +64,7 @@ def _manual_compression_reply_lines(summary: dict, compressor, focus_topic) -> l
lines.append(t(
"gateway.compress.aux_failed",
model=aux_fail_model,
error=(getattr(compressor, "_last_aux_model_failure_error", None) or "unknown error"),
))
error=(getattr(compressor, "_last_aux_model_failure_error", None) or "unknown error")))
return lines
@@ -78,11 +76,9 @@ def _compress_preview_reply(history, partial: bool, keep_last, focus_topic, agg_
pv_msgs = [
{"role": m.get("role"), "content": m.get("content")}
for m in history
if m.get("role") in {"user", "assistant"} and m.get("content")
]
if m.get("role") in {"user", "assistant"} and m.get("content")]
report = summarize_compress_preview(
pv_msgs, partial, keep_last, focus_topic, estimate_request_tokens_rough(pv_msgs)
)
pv_msgs, partial, keep_last, focus_topic, estimate_request_tokens_rough(pv_msgs))
lines = [f"🗜️ {line}" for line in report["lines"]]
if agg_note:
lines.append(agg_note)
@@ -136,31 +132,26 @@ class GatewaySessionCommandsMixin:
try:
await asyncio.wait_for(
self._run_in_executor_with_context(self._cleanup_agent_resources, _old_agent),
timeout=_RESET_CLEANUP_TIMEOUT_S,
)
timeout=_RESET_CLEANUP_TIMEOUT_S)
except asyncio.TimeoutError:
logger.warning(
"Agent resource cleanup for session %s exceeded %ss during /new reset; proceeding with "
"reset (the worker thread is left to finish on its own). (#35994)",
session_key, _RESET_CLEANUP_TIMEOUT_S,
)
session_key, _RESET_CLEANUP_TIMEOUT_S)
except Exception as cleanup_exc:
logger.warning(
"Agent resource cleanup for session %s failed during /new reset: %s (#35994)",
session_key, cleanup_exc,
)
session_key, cleanup_exc)
async def _fire_session_reset_hooks(
self, source: SessionSource, session_key: str, old_sid, new_sid
) -> None:
self, source: SessionSource, session_key: str, old_sid, new_sid) -> None:
"""Session-boundary hooks: plugin finalize (off-loop + bounded — trace exports can block
arbitrarily), then session:end and session:reset."""
platform_value = source.platform.value if source.platform else ""
with contextlib.suppress(Exception):
await self._finalize_session_off_loop(
session_id=old_sid, platform=platform_value, reason="new_session",
old_session_id=old_sid, new_session_id=new_sid,
)
old_session_id=old_sid, new_session_id=new_sid)
hook_payload = {"platform": platform_value, "user_id": source.user_id, "session_key": session_key}
await self.hooks.emit("session:end", dict(hook_payload))
await self.hooks.emit("session:reset", dict(hook_payload))
@@ -172,8 +163,7 @@ class GatewaySessionCommandsMixin:
_invoke_hook(
"on_session_reset", session_id=new_sid,
platform=source.platform.value if source.platform else "", reason="new_session",
old_session_id=old_sid, new_session_id=new_sid,
)
old_session_id=old_sid, new_session_id=new_sid)
except Exception:
pass
@@ -200,8 +190,7 @@ class GatewaySessionCommandsMixin:
interrupt_for_session(
session_key=session_key,
parent_session_id=str(getattr(old_entry, "session_id", "") or ""),
reason="session_reset",
)
reason="session_reset")
except Exception:
pass
_reset_process_scoped_tool_state()
@@ -209,8 +198,7 @@ class GatewaySessionCommandsMixin:
new_entry = await self.async_session_store.reset_session(session_key)
_old_sid = old_entry.session_id if old_entry else None
await self._fire_session_reset_hooks(
source, session_key, _old_sid, new_entry.session_id if new_entry else None
)
source, session_key, _old_sid, new_entry.session_id if new_entry else None)
# Scoped to the profile serving this source so a multiplexed /new banner reports the
# profile's model, not the base config's.
try:
@@ -234,8 +222,7 @@ class GatewaySessionCommandsMixin:
except Exception:
logger.debug("Failed to rebind Telegram topic after /new", exc_info=True)
self._invoke_session_reset_lifecycle_hook(
source, _old_sid, new_entry.session_id if new_entry else None
)
source, _old_sid, new_entry.session_id if new_entry else None)
try:
from hermes_cli.tips import get_random_tip
_tip_line = t("gateway.reset.tip", tip=get_random_tip())
@@ -277,8 +264,7 @@ class GatewaySessionCommandsMixin:
entries = getattr(self.session_store, "_entries", {}) or {}
return next(
(getattr(e, "origin", None) for e in entries.values() if getattr(e, "session_id", None) == session_id),
None,
)
None)
@staticmethod
def _same_matrix_room(current: SessionSource, origin: Optional[SessionSource]) -> bool:
@@ -289,8 +275,7 @@ class GatewaySessionCommandsMixin:
and origin.platform == Platform.MATRIX
and current.platform == Platform.MATRIX
and origin.chat_id == current.chat_id
and _sattr(current, "thread_id") == _sattr(origin, "thread_id")
)
and _sattr(current, "thread_id") == _sattr(origin, "thread_id"))
def _same_origin_chat(self, current: SessionSource, origin: Optional[SessionSource]) -> bool:
"""Platform-agnostic counterpart to ``_same_matrix_room``.
@@ -327,8 +312,7 @@ class GatewaySessionCommandsMixin:
build_session_key's isolation rules so the guards stay in lock-step with the key."""
return is_shared_multi_user_session(
source, group_sessions_per_user=getattr(self.config, "group_sessions_per_user", True),
thread_sessions_per_user=getattr(self.config, "thread_sessions_per_user", False),
)
thread_sessions_per_user=getattr(self.config, "thread_sessions_per_user", False))
def _resume_caller_is_admin(self, source: SessionSource) -> bool:
"""Whether *source* is an EXPLICITLY-configured admin (cross-origin /resume, /sessions).
@@ -381,8 +365,7 @@ class GatewaySessionCommandsMixin:
return bool(row_uid) and row_uid == caller_uid
async def _resume_target_allowed(
self, source: SessionSource, target_id: str, allow_override: bool = False
) -> bool:
self, source: SessionSource, target_id: str, allow_override: bool = False) -> bool:
"""Whether *source* may resume session *target_id* (IDOR guard for every adapter).
The live origin decides when the target is active; otherwise the DB row must PROVE
@@ -404,8 +387,7 @@ class GatewaySessionCommandsMixin:
return self._persisted_row_proves_owner(source, row)
async def _resume_row_visible(
self, source: SessionSource, row: dict, allow_all: bool
) -> bool:
self, source: SessionSource, row: dict, allow_all: bool) -> bool:
"""Whether a listing *row* belongs to the caller's origin (blocks cross-origin enumeration of
ids/previews); Matrix is room-scoped, ``--all`` needs a configured admin everywhere."""
if allow_all and self._resume_caller_is_admin(source):
@@ -425,16 +407,14 @@ class GatewaySessionCommandsMixin:
history_before_user_originated_turn,
retryable_user_text,
split_user_originated_turn,
user_originated_turn_view,
)
user_originated_turn_view)
source = event.source
session_entry = await self.async_session_store.get_or_create_session(source)
history = await self.async_session_store.load_transcript(session_entry.session_id)
last_user_idx = next(
(i for i in range(len(history) - 1, -1, -1) if user_originated_turn_view(history[i]) is not None),
None,
)
None)
if last_user_idx is None:
return t("gateway.retry.no_previous")
# Resolve text + scaffold-preserving prefix BEFORE any write; messaging retries cannot
@@ -452,8 +432,7 @@ class GatewaySessionCommandsMixin:
# on the same snapshot so a concurrent newer turn is never removed for stale text.
try:
rewind_result = await self.async_session_store.rewind_session(
session_entry.session_id, 1, require_retryable_composite=True,
)
session_entry.session_id, 1, require_retryable_composite=True)
except ValueError as exc:
return f"Cannot retry that message safely: {exc}"
if rewind_result is None:
@@ -461,14 +440,12 @@ class GatewaySessionCommandsMixin:
last_user_msg = rewind_result["target_text"]
# active_only preserves the active=0/compacted=1 archive left by in-place compaction.
elif not await self.async_session_store.rewrite_transcript(
session_entry.session_id, truncated, active_only=True, reject_active_turn_lease=True,
):
session_entry.session_id, truncated, active_only=True, reject_active_turn_lease=True):
return "Retry failed; transcript was not changed."
session_entry.last_prompt_tokens = 0 # transcript was truncated
retry_event = MessageEvent(
text=last_user_msg, message_type=MessageType.TEXT, source=source,
raw_message=event.raw_message, channel_prompt=event.channel_prompt,
)
raw_message=event.raw_message, channel_prompt=event.channel_prompt)
return await self._handle_message(retry_event)
async def _handle_undo_command(self, event: MessageEvent) -> str:
@@ -519,8 +496,7 @@ class GatewaySessionCommandsMixin:
return (
"🗜️ Nothing to compact: this session runs on the Codex app-server runtime, whose "
"context lives in a Codex-owned thread that only exists while the agent is active. "
"Send a message first, then /compress — or /reset to start fresh."
)
"Send a message first, then /compress — or /reset to start fresh.")
compressor = getattr(agent, "context_compressor", None)
count_before = getattr(compressor, "compression_count", 0)
try:
@@ -530,12 +506,10 @@ class GatewaySessionCommandsMixin:
if getattr(compressor, "compression_count", 0) > count_before:
return (
"🗜️ Codex app-server thread compacted (thread/compact). The transcript mirror is "
"unchanged by design — the app-server now carries the compacted context."
)
"unchanged by design — the app-server now carries the compacted context.")
return (
"⚠️ Codex app-server compaction did not complete — the thread is unchanged. Check the "
"app-server logs, retry /compress, or /reset for a clean session."
)
"app-server logs, retry /compress, or /reset for a clean session.")
async def _handle_compress_command_inner(self, event: MessageEvent) -> str:
"""Handle /compress -- manually compress conversation context; ``/compress <focus>`` tells
@@ -562,15 +536,13 @@ class GatewaySessionCommandsMixin:
return _compress_preview_reply(history, partial, keep_last, focus_topic, _agg_note)
try:
return await self._run_manual_compression(
source, session_entry, history, partial, keep_last, focus_topic
)
source, session_entry, history, partial, keep_last, focus_topic)
except Exception as e:
logger.warning("Manual compress failed: %s", e)
return t("gateway.compress.failed", error=e)
async def _run_manual_compression(
self, source, session_entry, history: list, partial: bool, keep_last, focus_topic
) -> str:
self, source, session_entry, history: list, partial: bool, keep_last, focus_topic) -> str:
"""Build a temporary agent, compress the transcript, persist, and describe the outcome."""
from agent.conversation_compression import finalize_context_engine_compression_notification
from agent.manual_compression_feedback import summarize_manual_compression
@@ -578,8 +550,7 @@ class GatewaySessionCommandsMixin:
from gateway.run import _platform_config_key
from hermes_cli.partial_compress import (
rejoin_compressed_head_and_tail,
split_history_for_partial_compress,
)
split_history_for_partial_compress)
session_key = self._session_key_for_source(source)
# Platform + stable gateway session key bind this agent (for external context engines) to
@@ -622,9 +593,7 @@ class GatewaySessionCommandsMixin:
compressed, _ = await self._run_in_executor_with_context(
lambda: tmp_agent._compress_context(
head, "", approx_tokens=approx_tokens, focus_topic=focus_topic, force=True,
defer_context_engine_notification=True,
)
)
defer_context_engine_notification=True))
# A held compression lock returns unchanged; say so instead of the misleading no-op text.
_lock_skipped = getattr(tmp_agent, "_compression_skipped_due_to_lock", None)
if _lock_skipped is True or isinstance(_lock_skipped, str):
@@ -636,8 +605,7 @@ class GatewaySessionCommandsMixin:
finalize_context_engine_compression_notification(tmp_agent, committed=True)
new_tokens = estimate_request_tokens_rough(compressed, system_prompt=_sys_prompt, tools=_tools)
summary = summarize_manual_compression(
msgs, compressed, approx_tokens, new_tokens, compression_state=compressor,
)
msgs, compressed, approx_tokens, new_tokens, compression_state=compressor)
finally:
finalize_context_engine_compression_notification(tmp_agent, committed=False)
self._evict_cached_agent(session_key) # next turn rebuilds the prompt from current files
@@ -663,20 +631,17 @@ class GatewaySessionCommandsMixin:
logger.warning(
"Manual compression could not restore the system prompt for session %s: %s. "
"Preserving an empty prompt so the live turn rebuilds it with its configured "
"providers.", session_id, exc, exc_info=True,
)
"providers.", session_id, exc, exc_info=True)
# compression.checkpoint_required needs the memory provider loaded so _compress_context()
# can write the pre-compression checkpoint; otherwise keep the fast path (no provider init).
_checkpoint_required = _is_truthy(
((_load_cfg() or {}).get("compression") or {}).get("checkpoint_required"),
default=False,
)
default=False)
tmp_agent = AIAgent(
**runtime_kwargs, model=model, max_iterations=4, quiet_mode=True,
skip_memory=not _checkpoint_required, enabled_toolsets=["memory"],
session_id=session_id, session_db=getattr(self._session_db, "_db", self._session_db),
)
session_id=session_id, session_db=getattr(self._session_db, "_db", self._session_db))
_seed_hygiene_system_prompt(tmp_agent, session_row)
# Real platform during construction (context engines bind correctly); afterwards a prompt
# rebuilt by compression is stamped as the provider-less fallback, stale for the next turn.
@@ -698,18 +663,15 @@ class GatewaySessionCommandsMixin:
if new_session_id != session_entry.session_id:
if not await self.async_session_store.rewrite_transcript(new_session_id, compressed):
raise RuntimeError(
f"failed to persist compressed transcript for session {new_session_id}"
)
f"failed to persist compressed transcript for session {new_session_id}")
session_entry.session_id = new_session_id
await self.async_session_store._save()
await asyncio.to_thread(
self._sync_telegram_topic_binding, source, session_entry, reason="compress-command",
)
self._sync_telegram_topic_binding, source, session_entry, reason="compress-command")
elif not getattr(tmp_agent, "_last_compaction_in_place", False):
logger.warning(
"Manual /compress: session rotation did not occur (session_id unchanged) and in-place "
"mode is off — preserving original transcript instead of overwriting it (#44794)."
)
"mode is off — preserving original transcript instead of overwriting it (#44794).")
await self.async_session_store.update_session(session_entry.session_key, last_prompt_tokens=0)
# ------------------------------------------------------------------------ /topic
@@ -756,8 +718,7 @@ class GatewaySessionCommandsMixin:
await self._session_db.enable_telegram_topic_mode(
chat_id=str(source.chat_id), user_id=str(source.user_id), profile_name=profile_name,
has_topics_enabled=capabilities.get("has_topics_enabled"),
allows_users_to_create_topics=capabilities.get("allows_users_to_create_topics"),
)
allows_users_to_create_topics=capabilities.get("allows_users_to_create_topics"))
except Exception as exc:
logger.exception("Failed to enable Telegram topic mode")
return t("gateway.topic.enable_failed", error=exc)
@@ -768,8 +729,7 @@ class GatewaySessionCommandsMixin:
try:
binding = await self._session_db.get_telegram_topic_binding(
chat_id=str(source.chat_id), thread_id=str(source.thread_id),
profile_name=profile_name,
)
profile_name=profile_name)
except Exception:
logger.debug("Failed to read Telegram topic binding", exc_info=True)
binding = None
@@ -782,8 +742,7 @@ class GatewaySessionCommandsMixin:
title = None
return t(
"gateway.topic.bound_status", label=title or t("gateway.topic.untitled_session"),
session_id=session_id,
)
session_id=session_id)
# ------------------------------------------------------------------ /save, /title
@@ -794,8 +753,7 @@ class GatewaySessionCommandsMixin:
SAVE_USAGE,
default_save_filename,
normalize_save_format,
render_session_for_save,
)
render_session_for_save)
parts = event.get_command_args().split()
redact = bool(parts) and parts[-1].lower() in ("redact", "--redact")
@@ -837,8 +795,7 @@ class GatewaySessionCommandsMixin:
return "Platform adapter not found to send the document."
await adapter.send_document(
chat_id=source.chat_id, file_path=temp_path, caption=f"Session export: {filename}",
file_name=filename,
)
file_name=filename)
return "Export complete."
except Exception as e:
logger.warning("Session /save failed: %s", e)
@@ -865,8 +822,7 @@ class GatewaySessionCommandsMixin:
session_id=session_id,
source=source.platform.value if source.platform else "unknown",
user_id=source.user_id, chat_id=source.chat_id, chat_type=source.chat_type,
thread_id=source.thread_id,
)
thread_id=source.thread_id)
title_arg = event.get_command_args().strip()
if not title_arg:
title = await self._session_db.get_session_title(session_id)
@@ -899,8 +855,7 @@ class GatewaySessionCommandsMixin:
widen = allow_all and self._resume_caller_is_admin(source)
sessions = await self._session_db.list_sessions_rich(
source=source.platform.value if source.platform else None,
session_key=None if widen else session_key, limit=10,
)
session_key=None if widen else session_key, limit=10)
titled = [s for s in sessions if s.get("title")][:10]
return [s for s in titled if await self._resume_row_visible(source, s, allow_all)]
@@ -942,8 +897,7 @@ class GatewaySessionCommandsMixin:
return t("gateway.resume.matrix_blocked_no_origin", name=name)
return t(
"gateway.resume.matrix_blocked_other_room",
room=target_origin.chat_name or target_origin.chat_id, name=name,
)
room=target_origin.chat_name or target_origin.chat_id, name=name)
if await self._resume_target_allowed(source, target_id, allow_override=(allow_all or allow_cross_room)):
return None
return t("gateway.resume.blocked_not_owner", name=name)
@@ -995,8 +949,7 @@ class GatewaySessionCommandsMixin:
msg_part = f" ({msg_count} message{'s' if msg_count != 1 else ''})" if msg_count else ""
return t(
"gateway.resume.matrix_cross_room_success", title=title,
room=source.chat_name or source.chat_id, msg_part=msg_part,
)
room=source.chat_name or source.chat_id, msg_part=msg_part)
if not msg_count:
return t("gateway.resume.resumed_no_count", title=title)
if msg_count == 1:
@@ -1036,13 +989,11 @@ class GatewaySessionCommandsMixin:
from hermes_cli.session_listing import (
format_gateway_session_listing,
parse_session_listing_args,
query_session_listing,
)
query_session_listing)
try:
include_all, include_unnamed, target, search_query = parse_session_listing_args(
event.get_command_args().strip()
)
event.get_command_args().strip())
except ValueError as exc:
return t("gateway.resume.parse_error", error=exc)
if search_query == "":
@@ -1070,8 +1021,7 @@ class GatewaySessionCommandsMixin:
search_query=search_query,
# Search filters in SQL: over-fetch so origin-invisible matches don't consume the page.
limit=50 if search_query else 10,
exclude_sources=["tool"],
)
exclude_sources=["tool"])
if not cross_origin:
rows = [row for row in rows if await self._resume_row_visible(source, row, allow_all=False)]
rows = rows[:10]
@@ -1124,8 +1074,7 @@ class GatewaySessionCommandsMixin:
chat_type=source.chat_type,
thread_id=source.thread_id,
origin_json=_branch_origin_json,
display_name=current_entry.display_name,
)
display_name=current_entry.display_name)
except Exception as e:
logger.error("Failed to create branch session: %s", e)
return t("gateway.branch.create_failed", error=e)
@@ -1133,8 +1082,7 @@ class GatewaySessionCommandsMixin:
# Chunked transactions; best-effort — a failed copy still yields a usable (partial) branch.
with contextlib.suppress(Exception):
await self._session_db.append_messages_batch(
new_session_id, [_branch_row(msg) for msg in history], chunk_rows=500,
)
new_session_id, [_branch_row(msg) for msg in history], chunk_rows=500)
with contextlib.suppress(Exception):
await self._session_db.set_session_title(new_session_id, branch_title)
new_entry = await self.async_session_store.switch_session(session_key, new_session_id)
+6 -12
View File
@@ -26,16 +26,13 @@ from gateway.platforms.base import _custom_unit_to_cp
from gateway.config import (
DEFAULT_STREAMING_EDIT_INTERVAL as _DEFAULT_STREAMING_EDIT_INTERVAL,
DEFAULT_STREAMING_BUFFER_THRESHOLD as _DEFAULT_STREAMING_BUFFER_THRESHOLD,
DEFAULT_STREAMING_CURSOR as _DEFAULT_STREAMING_CURSOR,
)
DEFAULT_STREAMING_CURSOR as _DEFAULT_STREAMING_CURSOR)
from gateway.response_filters import (
is_intentional_silence_response as _is_intentional_silence_response,
is_partial_silence_marker as _is_partial_silence_marker,
)
is_partial_silence_marker as _is_partial_silence_marker)
from gateway.stream_consumer_fences import ( # noqa: F401 (re-exported)
ensure_closed_code_fences,
escape_code_fences_for_display,
)
escape_code_fences_for_display)
from gateway.stream_consumer_transport import StreamTransportMixin
from gateway.stream_consumer_fallback import StreamFallbackMixin
from gateway.stream_consumer_think import StreamThinkFilterMixin
@@ -120,8 +117,7 @@ class GatewayStreamConsumer(StreamTransportMixin, StreamFallbackMixin, StreamThi
on_new_message: Optional[callable] = None,
on_before_finalize: Optional[Callable[[], Any]] = None,
initial_reply_to_id: Optional[str] = None,
run_still_current: Optional[Callable[[], bool]] = None,
):
run_still_current: Optional[Callable[[], bool]] = None):
self.adapter = adapter
self.chat_id = chat_id
self.cfg = config or StreamConsumerConfig()
@@ -781,8 +777,7 @@ class GatewayStreamConsumer(StreamTransportMixin, StreamFallbackMixin, StreamThi
# answer.
tick.update_visible = await self._send_or_edit(
display_text, finalize=tick.got_done or tick.got_segment_break,
is_turn_final=tick.got_done,
)
is_turn_final=tick.got_done)
self._last_edit_time = time.monotonic()
# Lines stay in _tool_progress_lines for the next compose.
self._tool_progress_active = False
@@ -885,8 +880,7 @@ class GatewayStreamConsumer(StreamTransportMixin, StreamFallbackMixin, StreamThi
if self._accumulated and self._message_id:
with contextlib.suppress(Exception):
best_effort_ok = bool(await self._send_or_edit(
self._accumulated, finalize=True, is_turn_final=False,
))
self._accumulated, finalize=True, is_turn_final=False))
elif self._message_id is None:
# Draft path keeps _message_id=None; seal in place (else the stream stays
# visibly live and the adapter keeps armed interception state).
+3 -6
View File
@@ -26,8 +26,7 @@ class StreamFallbackMixin:
try:
result = await self.adapter.send(
chat_id=self.chat_id, content=text, reply_to=reply_to_id,
metadata=self._metadata_for_send(final=final, expect_edits=not final),
)
metadata=self._metadata_for_send(final=final, expect_edits=not final))
if not (result.success and result.message_id):
self._edit_supported = False
return reply_to_id
@@ -216,8 +215,7 @@ class StreamFallbackMixin:
try:
result = await self._send_with_flood_retry(
content=final_text, reply_to=self._initial_reply_to_id,
retry_log="Flood control on empty fallback final send; retrying in %.1fs",
)
retry_log="Flood control on empty fallback final send; retrying in %.1fs")
except Exception as exc:
logger.debug("Empty fallback final send failed: %s", exc)
return "ambiguous" if self._send_failure_may_have_delivered(exc) else "failed"
@@ -323,8 +321,7 @@ class StreamFallbackMixin:
_needs_reply_anchor = _platform_name in ("buzz", "slack", "mattermost", "feishu")
result = await self.adapter.send(
chat_id=self.chat_id, content=text,
reply_to=self._initial_reply_to_id if _needs_reply_anchor else None, metadata=_md,
)
reply_to=self._initial_reply_to_id if _needs_reply_anchor else None, metadata=_md)
# Do NOT set _already_sent: commentary is interim, and the flag would
# suppress the real final after multiple tool calls.
if result.success:
+10 -20
View File
@@ -31,8 +31,7 @@ class StreamTransportMixin:
try:
params = inspect.signature(self.adapter.edit_message).parameters
if "metadata" in params or any(
param.kind is inspect.Parameter.VAR_KEYWORD for param in params.values()
):
param.kind is inspect.Parameter.VAR_KEYWORD for param in params.values()):
kwargs["metadata"] = self.metadata
except (TypeError, ValueError):
pass
@@ -43,8 +42,7 @@ class StreamTransportMixin:
bool; a raise logs ``fail_log`` at DEBUG (error formatted in, or the traceback when
``exc_info``) and reads as False."""
seed = self.adapter.send_stream_frame(
"", chat_id=self.chat_id, reply_to=self._initial_reply_to_id, turn_id=self._turn_id,
)
"", chat_id=self.chat_id, reply_to=self._initial_reply_to_id, turn_id=self._turn_id)
return await self._try_frame(seed, fail_log, exc_info=exc_info)
@staticmethod
@@ -63,8 +61,7 @@ class StreamTransportMixin:
"""One native-stream frame; every frame carries the same chat/reply/turn routing."""
return await self.adapter.send_stream_frame(
text, finalize=finalize, chat_id=self.chat_id, reply_to=self._initial_reply_to_id,
turn_id=self._turn_id,
)
turn_id=self._turn_id)
def _close_native_state(self) -> None:
"""Mark the native stream closed (next content re-seeds or falls back)."""
@@ -166,8 +163,7 @@ class StreamTransportMixin:
try:
result = await self.adapter.send_draft(
chat_id=self.chat_id, draft_id=self._draft_id, content=text,
metadata=self._draft_metadata(),
)
metadata=self._draft_metadata())
except Exception as e:
logger.debug("send_draft raised, disabling draft transport for this run: %s", e)
else:
@@ -191,8 +187,7 @@ class StreamTransportMixin:
try:
await self.adapter.abandon_open_draft(
self.chat_id, self._last_sent_text or self._clean_for_display(self._accumulated),
metadata=self._draft_metadata(),
)
metadata=self._draft_metadata())
except Exception as e:
logger.debug("abandon_open_draft failed (best-effort): %s", e)
@@ -255,8 +250,7 @@ class StreamTransportMixin:
stale_ids = self._stale_preview_ids()
try:
result = await self.adapter.send(
chat_id=self.chat_id, content=text, metadata=self._metadata_for_send(final=True),
)
chat_id=self.chat_id, content=text, metadata=self._metadata_for_send(final=True))
except Exception as e:
logger.debug("Fresh-final send failed, falling back to edit: %s", e)
return False
@@ -284,8 +278,7 @@ class StreamTransportMixin:
self._message_created_ts = None
async def _send_or_edit(
self, text: str, *, finalize: bool = False, is_turn_final: bool = True,
) -> bool:
self, text: str, *, finalize: bool = False, is_turn_final: bool = True) -> bool:
"""Send or edit the streaming message; True if delivered. ``finalize`` marks the
last edit. Transport order: native frame → draft frame → edit existing → first
send; a transport returns None to fall through to the next."""
@@ -413,8 +406,7 @@ class StreamTransportMixin:
"""First send, threaded to the user's message (correct topic/thread)."""
result = await self.adapter.send(
chat_id=self.chat_id, content=text, reply_to=self._initial_reply_to_id,
metadata=self._metadata_for_send(final=finalize, expect_edits=not finalize),
)
metadata=self._metadata_for_send(final=finalize, expect_edits=not finalize))
if not result.success:
self._edit_supported = False
return False
@@ -443,8 +435,7 @@ class StreamTransportMixin:
# CLASS (MagicMock auto-creates attrs) plus instance __dict__ (test doubles).
has_prefers_hook = (
hasattr(type(self.adapter), "prefers_fresh_final_streaming")
or "prefers_fresh_final_streaming" in getattr(self.adapter, "__dict__", {})
)
or "prefers_fresh_final_streaming" in getattr(self.adapter, "__dict__", {}))
prefers_fresh = self._adapter_prefers_fresh_final(text) # probed every edit (hook contract)
if finalize and (
prefers_fresh or (not has_prefers_hook and self._should_send_fresh_final())
@@ -518,8 +509,7 @@ class StreamTransportMixin:
logger.debug("Flood control on edit (strike %d/%d), backoff interval → %.1fs",
self._flood_strikes, self._MAX_FLOOD_STRIKES, self._current_edit_interval)
immediate_final_fallback = (
turn_final and getattr(self.adapter, "FALLBACK_ON_FINAL_EDIT_FLOOD", False) is True
)
turn_final and getattr(self.adapter, "FALLBACK_ON_FINAL_EDIT_FLOOD", False) is True)
if self._flood_strikes < self._MAX_FLOOD_STRIKES and not immediate_final_fallback:
self._last_edit_time = time.monotonic() # honor the new interval
return False