Skip to content
Open
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
53 changes: 44 additions & 9 deletions agent/turn_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -704,21 +704,56 @@ def _stamp_api_content_sidecar(
if _api_content is None or _api_content == _turn_user_msg.get("content"):
return
_turn_user_msg["api_content"] = _api_content
# In-place preflight compaction already inserted this turn's user row and the
# crash persist identity-skips compacted dicts, so backfill the stamp onto the row
# directly. Rotation mode flushes to the child session later.
if not (preflight_compressed and getattr(agent, "_last_compaction_in_place", False)):
# When this turn's user row was ALREADY materialized before the sidecar could be
# composed, the crash persist below skips the message (marker/identity) and the
# stamp would never reach the DB — the next turn then replays clean content and
# the request prefix diverges at this message. Two writers get there first:
# in-place preflight compaction (archive_and_compact runs before
# prefetch/pre_llm_call) and a close/early flush that raced the prologue on the
# CLI path (#102194). Both stamp ``_row_id`` on the live dict when they write it
# (``_insert_message_rows`` directly, ``sync_flushed_message_markers`` after the
# batch commit), so that id is at once the proof a row exists and the address to
# update — no positional guess, and no extra write on the normal path where the
# row does not exist yet.
#
# Do NOT widen this to an unconditional backfill: without a row id the store can
# only target the newest active user row, and a repeated user turn ("ok", "y",
# "continue") makes the PREVIOUS turn's row compare equal — this turn's bytes
# would overwrite its sidecar and be replayed as that turn forever.
#
# Rotation mode needs nothing here: its compacted copies flush to the child
# session after this stamp.
_row_id = _turn_user_msg.get("_row_id")
_has_valid_row_id = (
isinstance(_row_id, int)
and not isinstance(_row_id, bool)
and _row_id > 0
)
_in_place_compacted = preflight_compressed and bool(
getattr(agent, "_last_compaction_in_place", False)
)
if not (_has_valid_row_id or _in_place_compacted):
return
_db = getattr(agent, "_session_db", None)
if _db is not None:
try:
_db.set_latest_user_api_content(
agent.session_id, _turn_user_msg.get("content"), _api_content
)
if _has_valid_row_id:
_db.set_message_api_content(
agent.session_id,
_row_id,
_turn_user_msg.get("content"),
_api_content,
)
else:
# Compacted copy that carries no row id: fall back to the positional
# backfill, which is safe here because archive_and_compact just made
# this message the newest active user row.
_db.set_latest_user_api_content(
agent.session_id, _turn_user_msg.get("content"), _api_content
)
except Exception:
logger.warning(
"in-place compaction api_content backfill failed "
"for session=%s",
"api_content backfill failed for session=%s",
agent.session_id or "none",
exc_info=True,
)
Expand Down
21 changes: 21 additions & 0 deletions cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -2970,6 +2970,27 @@ def _claim_active_session(self, surface: str = "cli", *, stderr: bool = False) -
except Exception as exc:
logger.warning("Failed to claim active session slot: %s", exc)
return True
if message and getattr(message, "reason", None) == "SESSION_NOT_OWNED":
# Shared-brain deployments (gateway.per_session_exclusive: false) queue
# behind the live owner instead of being refused (#101279).
try:
from hermes_cli.active_sessions import per_session_exclusive, wait_for_session_ownership

if not per_session_exclusive(self.config) and wait_for_session_ownership(
session_id=self.session_id,
on_wait=lambda _e: self._console_print(
"[yellow]\N{HOURGLASS} Another Hermes process is using this session; "
"waiting for it to finish before starting your turn...[/yellow]"
),
):
lease, message = try_acquire_active_session(
session_id=self.session_id,
surface=surface,
config=self.config,
metadata={"live_session_id": str(self.session_id)},
)
except Exception:
logger.debug("Session ownership queueing unavailable; keeping refusal", exc_info=True)
if message:
print(message, file=sys.stderr) if stderr else self._console_print(f"[bold red]{message}[/]")
return False
Expand Down
92 changes: 92 additions & 0 deletions hermes_cli/active_sessions.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,40 @@ def resolve_max_concurrent_sessions(config: Any) -> Optional[int]:
return coerce_max_concurrent_sessions(raw, key=key)


def per_session_exclusive(config: Any) -> bool:
"""Resolve ``gateway.per_session_exclusive`` (default True) with top-level fallback.

True (default): a second live owner is REFUSED (``SESSION_NOT_OWNED``) — the solo-deployment
correctness default. False: shared-brain deployments may opt out, and a surfaced
``SESSION_NOT_OWNED`` refusal becomes a bounded WAIT for the current owner to finish (the
same serialization the messaging gateway already gives Telegram: the queued turn starts only
after the running one, with the full transcript). Invalid values warn once and keep the
safe default. See #101279.
"""
raw: Any = None
key = "gateway.per_session_exclusive"
if isinstance(config, dict):
gateway_cfg = config.get("gateway")
if isinstance(gateway_cfg, dict) and "per_session_exclusive" in gateway_cfg:
raw = gateway_cfg.get("per_session_exclusive")
elif "per_session_exclusive" in config:
raw = config.get("per_session_exclusive")
key = "per_session_exclusive"
else:
raw = getattr(config, "per_session_exclusive", None)
if raw is None:
gateway_cfg = getattr(config, "gateway", None)
raw = getattr(gateway_cfg, "per_session_exclusive", None)
if raw is None:
return True
if isinstance(raw, bool):
return raw
if isinstance(raw, str) and raw.strip().lower() in ("true", "false"):
return raw.strip().lower() == "true"
logger.warning("Ignoring invalid %s=%r (expected true/false); keeping exclusivity", key, raw)
return True


def format_age(seconds: float) -> str:
minutes = max(0, int(seconds // 60))
if minutes < 60:
Expand Down Expand Up @@ -526,6 +560,64 @@ def refuse(message: str, reason: str, log: str, *args) -> tuple[None, ActiveSess
return lease, None


def wait_for_session_ownership(
*, session_id: str, registry_home: str | Path | None = None,
wait_seconds: float = 1800.0, poll_seconds: float = 1.0,
should_abort=None, on_wait=None,
) -> bool:
"""Block until no OTHER live lease holds *session_id* (or the wait bounds expire).

The queueing half of ``gateway.per_session_exclusive: false`` (#101279): after a
``SESSION_NOT_OWNED`` refusal, the caller waits for the current owner to finish and
re-acquires — the serialization the messaging gateway already gives Telegram, instead
of pushing a refusal onto the second user. Returns True when the session is free (or
was free all along); False when the wait timed out or was aborted. Never raises: a
registry that cannot be read mid-wait returns False so the caller falls back to the
refusal path.

*should_abort* (optional): called each poll; a truthy return cancels the wait.
*on_wait* (optional): called once with the elapsed seconds when the wait proves
non-trivial (>0 polls), for a status line like the gateway's
"Another Hermes process is using this session; waiting...".
"""
key = str(session_id or "")
deadline = time.monotonic() + max(0.0, float(wait_seconds))
poll = max(0.05, float(poll_seconds))
notified = False
state_path = _state_path(registry_home)
lock_path = _lock_path(registry_home)
while True:
try:
with _FileLock(lock_path):
loaded = _read_live_entries(
state_path, track_liveness=False,
warn="Active-session registry is unavailable while waiting for ownership",
)
if loaded is None:
return False
_, entries = loaded
if not any(
str(existing.get("session_id") or "") == key for existing in entries
):
return True
except Exception:
logger.warning("Ownership wait registry check failed", exc_info=True)
return False
if should_abort is not None:
try:
if should_abort():
return False
except Exception:
logger.debug("Ownership wait abort probe failed", exc_info=True)
if time.monotonic() >= deadline:
return False
if on_wait is not None and not notified:
notified = True
with suppress(Exception):
on_wait(0.0)
time.sleep(poll)


def release_active_session(lease: ActiveSessionLease) -> None:
# Prefer the registry the lease was acquired against: the caller may be
# running under a profile HERMES_HOME override.
Expand Down
31 changes: 30 additions & 1 deletion hermes_state_messages.py
Original file line number Diff line number Diff line change
Expand Up @@ -576,13 +576,42 @@ def _message_column_names(self, conn) -> List[str]:
def set_latest_user_api_content(self, session_id: str, content: Any, api_content: str) -> int:
"""Backfill the ``api_content`` sidecar onto the newest ACTIVE user row (0/1 rows). Preflight compaction
inserts that row BEFORE the sidecar exists and the later persist identity-skips compacted dicts;
without this a reload reopens the prompt-cache divergence. ``content`` match guards a racing rewrite."""
without this a reload reopens the prompt-cache divergence. ``content`` match guards a racing rewrite.

POSITIONAL, and only safe when the caller already knows the newest active user row IS the message it
stamped. The content match is NOT sufficient on its own: repeated identical user turns ("ok", "y",
"continue") make an OLDER row compare equal, so calling this before the current turn's row exists
overwrites the previous turn's sidecar with this turn's bytes — durable wrong-bytes replay, a worse
cache break than the missing sidecar. When the caller holds the durable row id (``_row_id``, synced
onto the live dict by ``sync_flushed_message_markers`` and stamped by ``_insert_message_rows``), use
:meth:`set_message_api_content` instead — it addresses the exact row and cannot land on a neighbour."""
return self._write_rowcount(
"UPDATE messages SET api_content = ? WHERE id = (SELECT id FROM messages "
"WHERE session_id = ? AND role = 'user' AND active = 1 ORDER BY id DESC LIMIT 1"
") AND content IS ?",
(_scrub_surrogates(api_content), session_id, self._encode_content(content)))

def set_message_api_content(self, session_id: str, row_id: int, content: Any, api_content: str) -> int:
"""Backfill the ``api_content`` sidecar onto ONE known durable row.

Row-addressed counterpart to :meth:`set_latest_user_api_content`: the caller passes the ``_row_id``
the write path stamped on the live message dict, so the update cannot drift onto a neighbouring row
that merely carries the same text.

Used by the turn prologue whenever the current turn's user row was already materialized before the
sidecar could be composed (in-place preflight compaction, a close/early flush that raced the
prologue — #102194). The crash persist then marker-skips that message, so this is the only way the
stamped bytes reach the store.

``active = 1`` and the ``content`` match stay as defensive guards: a row the compaction archived, or
one a racing rewrite changed, is left untouched."""
if not session_id or isinstance(row_id, bool) or not isinstance(row_id, int) or row_id <= 0:
return 0
return self._write_rowcount(
"UPDATE messages SET api_content = ? "
"WHERE id = ? AND session_id = ? AND role = 'user' AND active = 1 AND content IS ?",
(_scrub_surrogates(api_content), row_id, session_id, self._encode_content(content)))

def _dedupe_display_generations(self, rows):
"""Collapse compaction generations so each logical message appears once (the protected tail is copied
into each generation: same role/content/timestamp, different ``active``/id); prefer the live row, then
Expand Down
Loading
Loading