From f2f0be898f29317b449784907ceec8b19f649d8a Mon Sep 17 00:00:00 2001 From: KoNit-K Date: Thu, 17 Sep 2026 12:30:48 +0800 Subject: [PATCH] fix(state): support legacy Telegram topic schemas profile purge/rekey gate their profile_name SQL on the column actually being present, so supported v1/v2 topic tables (never migrated because /topic was never run) no longer abort identity settlement with "no such column: profile_name". Binding cleanup/rekey still runs by exact agent:: session_key namespace on those shapes. Fixes #113757 Co-authored-by: Yagna Vudathu --- hermes_state_gateway.py | 35 ++++++++--- .../hermes_state/test_purge_profile_state.py | 61 +++++++++++++++++++ .../hermes_state/test_rekey_profile_state.py | 57 +++++++++++++++++ 3 files changed, 144 insertions(+), 9 deletions(-) diff --git a/hermes_state_gateway.py b/hermes_state_gateway.py index ac7bd5a1a1..2166bd0724 100644 --- a/hermes_state_gateway.py +++ b/hermes_state_gateway.py @@ -558,6 +558,11 @@ class SessionGatewayMixin: def _do(conn): existing = {row[0] for row in conn.execute( "SELECT name FROM sqlite_master WHERE type='table'").fetchall()} + topic_columns = { + table: {row[1] for row in conn.execute(f"PRAGMA table_info('{table}')")} + for table in ("telegram_dm_topic_mode", "telegram_dm_topic_bindings") + if table in existing + } collision = conn.execute( "SELECT old.scope, ? || substr(old.session_key, ?) " "FROM gateway_routing AS old JOIN gateway_routing AS target " @@ -573,7 +578,7 @@ class SessionGatewayMixin: ("telegram_dm_topic_mode", ("chat_id",)), ("telegram_dm_topic_bindings", ("chat_id", "thread_id")), ): - if table not in existing: + if "profile_name" not in topic_columns.get(table, set()): continue equality = " AND ".join( f"target.{column} = old.{column}" for column in columns) @@ -616,11 +621,11 @@ class SessionGatewayMixin: "WHERE substr(session_key, 1, ?) = ?", (new_ns, ns_len + 1, ns_len, old_ns)).rowcount for table in ("telegram_dm_topic_mode", "telegram_dm_topic_bindings"): - if table in existing: + if "profile_name" in topic_columns.get(table, set()): counts[f"{table}_profile_name"] = conn.execute( f"UPDATE {table} SET profile_name = ? WHERE profile_name = ?", (new, old)).rowcount - if "telegram_dm_topic_bindings" in existing: + if "session_key" in topic_columns.get("telegram_dm_topic_bindings", set()): counts["telegram_dm_topic_bindings_session_key"] = conn.execute( "UPDATE telegram_dm_topic_bindings " "SET session_key = ? || substr(session_key, ?) " @@ -686,6 +691,11 @@ class SessionGatewayMixin: def _do(conn): existing = {row[0] for row in conn.execute( "SELECT name FROM sqlite_master WHERE type='table'").fetchall()} + topic_columns = { + table: {row[1] for row in conn.execute(f"PRAGMA table_info('{table}')")} + for table in ("telegram_dm_topic_mode", "telegram_dm_topic_bindings") + if table in existing + } if "gateway_routing" in existing: counts["gateway_routing"] = conn.execute( "DELETE FROM gateway_routing WHERE substr(session_key, 1, ?) = ?", @@ -702,16 +712,23 @@ class SessionGatewayMixin: "WHERE (adapter_profile = ? OR substr(session_key, 1, ?) = ?) " "AND state NOT IN ('delivered', 'abandoned')", (time.time(), name, ns_len, ns)).rowcount - if "telegram_dm_topic_mode" in existing: + if "profile_name" in topic_columns.get("telegram_dm_topic_mode", set()): counts["telegram_dm_topic_mode"] = conn.execute( "DELETE FROM telegram_dm_topic_mode WHERE profile_name = ?", (name,)).rowcount - if "telegram_dm_topic_bindings" in existing: + binding_columns = topic_columns.get("telegram_dm_topic_bindings", set()) + if "session_key" in binding_columns: # A rename rewrites a binding's session_key namespace as well as its profile_name # (:meth:`rekey_profile_state`), so matching on one alone leaves the other behind. - counts["telegram_dm_topic_bindings"] = conn.execute( - "DELETE FROM telegram_dm_topic_bindings " - "WHERE profile_name = ? OR substr(session_key, 1, ?) = ?", - (name, ns_len, ns)).rowcount + if "profile_name" in binding_columns: + counts["telegram_dm_topic_bindings"] = conn.execute( + "DELETE FROM telegram_dm_topic_bindings " + "WHERE profile_name = ? OR substr(session_key, 1, ?) = ?", + (name, ns_len, ns)).rowcount + else: + counts["telegram_dm_topic_bindings"] = conn.execute( + "DELETE FROM telegram_dm_topic_bindings " + "WHERE substr(session_key, 1, ?) = ?", + (ns_len, ns)).rowcount self._execute_write(_do) return counts diff --git a/tests/hermes_state/test_purge_profile_state.py b/tests/hermes_state/test_purge_profile_state.py index dc3e2a9f98..5e51304a70 100644 --- a/tests/hermes_state/test_purge_profile_state.py +++ b/tests/hermes_state/test_purge_profile_state.py @@ -21,6 +21,29 @@ def db(tmp_path): database.close() +def _create_legacy_v2_topic_tables(db): + """Create the supported pre-profile-name topic schema without migrating it.""" + db._write_sql(""" + CREATE TABLE telegram_dm_topic_mode ( + chat_id TEXT PRIMARY KEY, user_id TEXT NOT NULL, + enabled INTEGER NOT NULL DEFAULT 1, + activated_at REAL NOT NULL, updated_at REAL NOT NULL, + has_topics_enabled INTEGER, allows_users_to_create_topics INTEGER, + capability_checked_at REAL, intro_message_id TEXT, pinned_message_id TEXT + ) + """) + db._write_sql(""" + CREATE TABLE telegram_dm_topic_bindings ( + chat_id TEXT NOT NULL, thread_id TEXT NOT NULL, user_id TEXT NOT NULL, + session_key TEXT NOT NULL, + session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, + managed_mode TEXT NOT NULL DEFAULT 'auto', + linked_at REAL NOT NULL, updated_at REAL NOT NULL, + PRIMARY KEY (chat_id, thread_id) + ) + """) + + def _write_obligation(db, monkeypatch, obligation_id, session_key, chat_id, profile, state="pending"): """delivery_obligations is created lazily by the delivery ledger against the same state.db.""" @@ -150,3 +173,41 @@ class TestPurgeProfileState: routing = db.load_gateway_routing_entries(scope="/root/sessions") assert "agent:foo_bar:feishu:dm:chatA" not in routing assert "agent:fooXbar:feishu:dm:chatB" in routing + + def test_purges_legacy_v2_topic_binding_by_session_key(self, db): + """Purge is safe on old tables and leaves unrelated namespace rows intact.""" + db.create_session( + "sess_gone", "telegram", session_key="agent:gone:telegram:dm:chatA", + profile_name="gone", chat_id="chatA", chat_type="dm") + db.create_session( + "sess_keep", "telegram", session_key="agent:keepme:telegram:dm:chatB", + profile_name="keepme", chat_id="chatB", chat_type="dm") + _create_legacy_v2_topic_tables(db) + db._write_sql( + "INSERT INTO telegram_dm_topic_bindings " + "(chat_id, thread_id, user_id, session_key, session_id, linked_at, updated_at) " + "VALUES (?, ?, ?, ?, ?, 1, 1)", + ("chatA", "threadA", "userA", "agent:gone:telegram:dm:chatA", "sess_gone")) + db._write_sql( + "INSERT INTO telegram_dm_topic_bindings " + "(chat_id, thread_id, user_id, session_key, session_id, linked_at, updated_at) " + "VALUES (?, ?, ?, ?, ?, 1, 1)", + ("chatB", "threadB", "userB", "agent:keepme:telegram:dm:chatB", "sess_keep")) + + counts = db.purge_profile_state("gone") + + assert counts["telegram_dm_topic_bindings"] == 1 + assert db._read_one( + "SELECT COUNT(*) AS n FROM telegram_dm_topic_bindings WHERE chat_id = ?", ("chatA",) + )["n"] == 0 + assert db._read_one( + "SELECT session_key FROM telegram_dm_topic_bindings WHERE chat_id = ?", ("chatB",) + )["session_key"] == "agent:keepme:telegram:dm:chatB" + + def test_purges_zero_legacy_v2_rows_without_profile_column(self, db): + """A delete for an absent profile must not query profile_name on a v2 schema.""" + _create_legacy_v2_topic_tables(db) + + counts = db.purge_profile_state("gone") + + assert counts["telegram_dm_topic_bindings"] == 0 diff --git a/tests/hermes_state/test_rekey_profile_state.py b/tests/hermes_state/test_rekey_profile_state.py index 92ba91a1fc..3a7e856b59 100644 --- a/tests/hermes_state/test_rekey_profile_state.py +++ b/tests/hermes_state/test_rekey_profile_state.py @@ -21,6 +21,29 @@ def db(tmp_path): database.close() +def _create_legacy_v2_topic_tables(db): + """Create the supported pre-profile-name topic schema without migrating it.""" + db._write_sql(""" + CREATE TABLE telegram_dm_topic_mode ( + chat_id TEXT PRIMARY KEY, user_id TEXT NOT NULL, + enabled INTEGER NOT NULL DEFAULT 1, + activated_at REAL NOT NULL, updated_at REAL NOT NULL, + has_topics_enabled INTEGER, allows_users_to_create_topics INTEGER, + capability_checked_at REAL, intro_message_id TEXT, pinned_message_id TEXT + ) + """) + db._write_sql(""" + CREATE TABLE telegram_dm_topic_bindings ( + chat_id TEXT NOT NULL, thread_id TEXT NOT NULL, user_id TEXT NOT NULL, + session_key TEXT NOT NULL, + session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, + managed_mode TEXT NOT NULL DEFAULT 'auto', + linked_at REAL NOT NULL, updated_at REAL NOT NULL, + PRIMARY KEY (chat_id, thread_id) + ) + """) + + class TestRekeyProfileState: def test_rekeys_session_key_namespace_and_profile_columns(self, db): # A session owned by the old profile, keyed in its namespace. @@ -137,3 +160,37 @@ class TestRekeyProfileState: "WHERE chat_id = ? AND thread_id = ?", ("chatA", "threadA")) assert binding["profile_name"] == "newname" assert binding["session_key"] == "agent:newname:telegram:dm:chatA" + + def test_rekeys_legacy_v2_topic_binding_by_session_key(self, db): + """Legacy v2 tables have no profile_name, but bindings retain profile namespaces.""" + db.create_session( + "sess_old", "telegram", session_key="agent:oldname:telegram:dm:chatA", + profile_name="oldname", chat_id="chatA", chat_type="dm") + db.create_session( + "sess_keep", "telegram", session_key="agent:keepme:telegram:dm:chatB", + profile_name="keepme", chat_id="chatB", chat_type="dm") + _create_legacy_v2_topic_tables(db) + db._write_sql( + "INSERT INTO telegram_dm_topic_mode " + "(chat_id, user_id, enabled, activated_at, updated_at) VALUES (?, ?, 1, 1, 1)", + ("chatA", "userA")) + db._write_sql( + "INSERT INTO telegram_dm_topic_bindings " + "(chat_id, thread_id, user_id, session_key, session_id, linked_at, updated_at) " + "VALUES (?, ?, ?, ?, ?, 1, 1)", + ("chatA", "threadA", "userA", "agent:oldname:telegram:dm:chatA", "sess_old")) + db._write_sql( + "INSERT INTO telegram_dm_topic_bindings " + "(chat_id, thread_id, user_id, session_key, session_id, linked_at, updated_at) " + "VALUES (?, ?, ?, ?, ?, 1, 1)", + ("chatB", "threadB", "userB", "agent:keepme:telegram:dm:chatB", "sess_keep")) + + counts = db.rekey_profile_state("oldname", "newname") + + assert counts["telegram_dm_topic_bindings_session_key"] == 1 + assert db._read_one( + "SELECT session_key FROM telegram_dm_topic_bindings WHERE chat_id = ?", ("chatA",) + )["session_key"] == "agent:newname:telegram:dm:chatA" + assert db._read_one( + "SELECT session_key FROM telegram_dm_topic_bindings WHERE chat_id = ?", ("chatB",) + )["session_key"] == "agent:keepme:telegram:dm:chatB"