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:<name>: session_key namespace on those shapes. Fixes #113757 Co-authored-by: Yagna Vudathu <yagnavudathu@gmail.com>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user