diff --git a/hermes_state_compression.py b/hermes_state_compression.py index eb9dfb5f49..9786f97a28 100644 --- a/hermes_state_compression.py +++ b/hermes_state_compression.py @@ -102,8 +102,7 @@ class SessionCompressionMixin: return False updated = conn.execute( "UPDATE sessions SET ended_at = NULL, end_reason = NULL " - "WHERE id = ? AND ended_at IS NOT NULL " - "AND end_reason = 'compression'", + "WHERE id = ? AND ended_at IS NOT NULL AND end_reason = 'compression'", (session_id,), ) # rowcount==1 is guaranteed by the parent SELECT in this same txn. A False @@ -115,7 +114,10 @@ class SessionCompressionMixin: def _publish_child_session_row(self, conn, parent, *, parent_session_id, child_session_id, source, model, model_config, system_prompt, cwd, profile_name) -> None: - """INSERT the compression child's ``sessions`` row copied from *parent*.""" + """INSERT the compression child's ``sessions`` row copied from *parent*. Same contract + as _insert_session_row's compression-fork backfill: the child stays on the parent's + profile and keeps gateway routing/origin columns; no owner on either side -> this + store's profile.""" system_prompt_hash = self._store_system_prompt(conn, system_prompt) conn.execute( """INSERT INTO sessions ( @@ -126,27 +128,12 @@ class SessionCompressionMixin: thread_id, display_name, origin_json, started_at ) VALUES (?, ?, ?, ?, NULL, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", ( - child_session_id, - source, - model, - json.dumps(model_config) if model_config else None, - system_prompt_hash, - parent_session_id, - cwd or parent["cwd"], - parent["git_branch"], + child_session_id, source, model, json.dumps(model_config) if model_config else None, + system_prompt_hash, parent_session_id, cwd or parent["cwd"], parent["git_branch"], parent["git_repo_root"], - # Same contract as _insert_session_row's compression-fork backfill: the - # child stays on the parent's profile and keeps gateway routing/origin - # columns; no owner on either side -> this store's profile. profile_name or parent["profile_name"] or self._own_profile_name(), - parent["user_id"], - parent["session_key"], - parent["chat_id"], - parent["chat_type"], - parent["thread_id"], - parent["display_name"], - parent["origin_json"], - time.time(), + parent["user_id"], parent["session_key"], parent["chat_id"], parent["chat_type"], + parent["thread_id"], parent["display_name"], parent["origin_json"], time.time(), ), ) @@ -209,8 +196,7 @@ class SessionCompressionMixin: # Deliberate boundaries still fail closed. if is_automatic_end_reason(parent["end_reason"]): conn.execute( - "UPDATE sessions SET ended_at = NULL, end_reason = NULL " - "WHERE id = ?", + "UPDATE sessions SET ended_at = NULL, end_reason = NULL WHERE id = ?", (parent_session_id,), ) else: @@ -419,8 +405,7 @@ class SessionCompressionMixin: expires_at = time.time() + ttl_seconds try: return self._write_rowcount( - "UPDATE compression_locks SET expires_at = ? " - "WHERE session_id = ? AND holder = ?", + "UPDATE compression_locks SET expires_at = ? WHERE session_id = ? AND holder = ?", (expires_at, session_id, holder), ) > 0 except sqlite3.Error as exc: @@ -448,15 +433,13 @@ class SessionCompressionMixin: current_holder, current_expires_at = row[0], row[1] if current_expires_at < now or _compression_lock_holder_process_is_dead(current_holder): conn.execute( - "DELETE FROM compression_locks " - "WHERE session_id = ? AND holder = ?", + "DELETE FROM compression_locks WHERE session_id = ? AND holder = ?", (session_id, current_holder), ) reclaimed_holder = current_holder conn.execute( "INSERT OR IGNORE INTO compression_locks " - "(session_id, holder, acquired_at, expires_at) " - "VALUES (?, ?, ?, ?)", + "(session_id, holder, acquired_at, expires_at) VALUES (?, ?, ?, ?)", (session_id, holder, now, expires_at), ) row = conn.execute("SELECT holder FROM compression_locks WHERE session_id = ?", (session_id,)).fetchone() @@ -480,8 +463,7 @@ class SessionCompressionMixin: return self._write_sql_logged( "release_compression_lock", session_id, - "DELETE FROM compression_locks " - "WHERE session_id = ? AND holder = ?", + "DELETE FROM compression_locks WHERE session_id = ? AND holder = ?", (session_id, holder), ) @@ -538,22 +520,19 @@ class SessionCompressionMixin: def _do(conn): conversation_id = self._session_turn_lease_key_on_conn(conn, session_id) row = conn.execute( - "SELECT holder, expires_at FROM session_turn_leases " - "WHERE conversation_id = ?", + "SELECT holder, expires_at FROM session_turn_leases WHERE conversation_id = ?", (conversation_id,), ).fetchone() if row is not None: current_holder = row["holder"] if float(row["expires_at"]) <= now or _compression_lock_holder_process_is_dead(current_holder): conn.execute( - "DELETE FROM session_turn_leases " - "WHERE conversation_id = ? AND holder = ?", + "DELETE FROM session_turn_leases WHERE conversation_id = ? AND holder = ?", (conversation_id, current_holder), ) conn.execute( "INSERT OR IGNORE INTO session_turn_leases " - "(conversation_id, holder, acquired_at, expires_at) " - "VALUES (?, ?, ?, ?)", + "(conversation_id, holder, acquired_at, expires_at) VALUES (?, ?, ?, ?)", (conversation_id, holder, now, expires_at), ) owner = conn.execute( @@ -637,8 +616,7 @@ class SessionCompressionMixin: def _do(conn): conversation_id = self._session_turn_lease_key_on_conn(conn, session_id) conn.execute( - "DELETE FROM session_turn_leases " - "WHERE conversation_id = ? AND holder = ?", + "DELETE FROM session_turn_leases WHERE conversation_id = ? AND holder = ?", (conversation_id, holder), ) @@ -649,8 +627,7 @@ class SessionCompressionMixin: if not session_id: return None row = self._read_one( - "SELECT holder FROM compression_locks " - "WHERE session_id = ? AND expires_at >= ?", + "SELECT holder FROM compression_locks WHERE session_id = ? AND expires_at >= ?", (session_id, time.time()), ) return None if row is None else row[0] diff --git a/hermes_state_messages.py b/hermes_state_messages.py index 3dbe8f3555..a318c4487a 100644 --- a/hermes_state_messages.py +++ b/hermes_state_messages.py @@ -91,6 +91,12 @@ def _ended_by_compression(row) -> bool: return row is not None and row["ended_at"] is not None and row["end_reason"] == "compression" +def _stale_holder(row, now: float) -> bool: + """A lock/lease row whose holder is expired or a provably dead local process.""" + from hermes_state import _compression_lock_holder_process_is_dead + return float(row["expires_at"]) <= now or _compression_lock_holder_process_is_dead(row["holder"]) + + class SessionMessagesMixin: """Message append/replace/rewind, reactions, resume conversations, replay dedupe.""" @@ -216,17 +222,13 @@ class SessionMessagesMixin: ``reject_active_compression_lock`` / ``reject_active_turn_lease`` so a compressor that captured its watermark cannot resurrect the removed turn. """ - from hermes_state import CompressionSessionClosedError, SessionCompressionInProgressError, SessionTurnLeaseLostError, _compression_lock_holder_process_is_dead + from hermes_state import CompressionSessionClosedError, SessionCompressionInProgressError, SessionTurnLeaseLostError if reject_active_compression_lock: active_lock = conn.execute(_COMPRESSION_LOCK_ROW_SQL, (session_id,)).fetchone() if active_lock is not None: - current_holder = active_lock["holder"] - if ( - float(active_lock["expires_at"]) <= time.time() - or _compression_lock_holder_process_is_dead(current_holder) - ): - conn.execute(_DELETE_COMPRESSION_LOCK_SQL, (session_id, current_holder)) - elif current_holder != compression_lock_holder: + if _stale_holder(active_lock, time.time()): + conn.execute(_DELETE_COMPRESSION_LOCK_SQL, (session_id, active_lock["holder"])) + elif active_lock["holder"] != compression_lock_holder: raise SessionCompressionInProgressError( f"Session {session_id!r} is being compressed by another writer" ) @@ -249,21 +251,16 @@ class SessionMessagesMixin: (now + max(0.1, float(turn_lease_ttl_seconds)), conversation_id, turn_lease_holder), ) elif lease is not None: - current_holder = lease["holder"] - if ( - float(lease["expires_at"]) <= now - or _compression_lock_holder_process_is_dead(current_holder) - ): - # Same reclaim rule as acquisition (expired or provably dead owner); - # deleting here also fences a stale late flush after the mutation. - conn.execute( - "DELETE FROM session_turn_leases WHERE conversation_id = ? AND holder = ?", - (conversation_id, current_holder), - ) - else: + if not _stale_holder(lease, now): raise SessionTurnLeaseLostError( f"Session has an active turn lease; refusing transcript mutation for {session_id!r}" ) + # Same reclaim rule as acquisition (expired or provably dead owner); + # deleting here also fences a stale late flush after the mutation. + conn.execute( + "DELETE FROM session_turn_leases WHERE conversation_id = ? AND holder = ?", + (conversation_id, lease["holder"]), + ) session = conn.execute(_ENDED_BY_COMPRESSION_SQL, (session_id,)).fetchone() if _ended_by_compression(session) and not allow_closed_compression_parent: raise CompressionSessionClosedError(session_id) @@ -519,8 +516,7 @@ class SessionMessagesMixin: def _do(conn): rows = conn.execute( "SELECT id, role, content, display_metadata FROM messages " - "WHERE session_id = ? AND active = 1 AND display_metadata IS NOT NULL " - "ORDER BY id", + "WHERE session_id = ? AND active = 1 AND display_metadata IS NOT NULL ORDER BY id", (session_id,), ).fetchall() pending = [] @@ -758,8 +754,7 @@ class SessionMessagesMixin: tail_ids, tail_tool_calls = self._tail_rows_after_watermark( conn, "SELECT id, tool_calls FROM messages " - "WHERE session_id = ? AND active = 1 AND id > ? " - "ORDER BY id", + "WHERE session_id = ? AND active = 1 AND id > ? ORDER BY id", (session_id, int(watermark)), ) # Rewind targets sit AT/BELOW the watermark (the compressor only saw rows up @@ -768,8 +763,7 @@ class SessionMessagesMixin: if tail_count > 0: if watermark is not None: tail_rows = conn.execute( - "SELECT id FROM messages " - "WHERE session_id = ? AND active = 1 AND id <= ? " + "SELECT id FROM messages WHERE session_id = ? AND active = 1 AND id <= ? " "ORDER BY id DESC LIMIT ?", (session_id, int(watermark), int(tail_count)), ).fetchall() @@ -840,10 +834,8 @@ class SessionMessagesMixin: """ from hermes_state import _scrub_surrogates 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" + "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)), ) @@ -991,15 +983,11 @@ class SessionMessagesMixin: if not anchor_exists: return {"window": [], "messages_before": 0, "messages_after": 0} before_rows = conn.execute( - "SELECT * FROM messages " - "WHERE session_id = ? AND id <= ? " - "ORDER BY id DESC LIMIT ?", + "SELECT * FROM messages WHERE session_id = ? AND id <= ? ORDER BY id DESC LIMIT ?", (session_id, around_message_id, window + 1), ).fetchall() after_rows = conn.execute( - "SELECT * FROM messages " - "WHERE session_id = ? AND id > ? " - "ORDER BY id ASC LIMIT ?", + "SELECT * FROM messages WHERE session_id = ? AND id > ? ORDER BY id ASC LIMIT ?", (session_id, around_message_id, window), ).fetchall() rows = list(reversed(before_rows)) + list(after_rows) @@ -1046,8 +1034,7 @@ class SessionMessagesMixin: best = current try: child_row = conn.execute( - "SELECT id FROM sessions AS child " - "WHERE child.parent_session_id = ? " + "SELECT id FROM sessions AS child WHERE child.parent_session_id = ? " " AND json_extract(COALESCE(child.model_config, '{}'), '$._branched_from') IS NULL " " AND json_extract(COALESCE(child.model_config, '{}'), '$._delegate_from') IS NULL " " AND json_extract(COALESCE(child.model_config, '{}'), '$._reset_from') IS NULL " diff --git a/hermes_state_titles.py b/hermes_state_titles.py index 9404fa74ad..c6ac1dfdc6 100644 --- a/hermes_state_titles.py +++ b/hermes_state_titles.py @@ -166,18 +166,15 @@ class SessionTitlesMixin: if source not in self._TITLE_SOURCE_RANK: raise ValueError(f"invalid title source: {source!r}") return self._write_rowcount( - "UPDATE sessions SET title_source = ? " - "WHERE id = ? AND title IS NOT NULL", + "UPDATE sessions SET title_source = ? WHERE id = ? AND title IS NOT NULL", (source, session_id), ) > 0 def get_session_by_title(self, title: str) -> Optional[Dict[str, Any]]: """Look up a session by exact title. Returns session dict or None.""" row = self._read_one( - "SELECT s.*, " - "COALESCE(sp.prompt, s.system_prompt) AS _system_prompt_resolved " - "FROM sessions s " - "LEFT JOIN system_prompts sp ON sp.hash = s.system_prompt_hash " + "SELECT s.*, COALESCE(sp.prompt, s.system_prompt) AS _system_prompt_resolved " + "FROM sessions s LEFT JOIN system_prompts sp ON sp.hash = s.system_prompt_hash " "WHERE s.title = ?", (title,), )