From fa827f596e195a18b962bdf862eff8ca86dec3b4 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 16:46:00 -0700 Subject: [PATCH] refactor(gateway/hosted_rooms): reflow multi-line SQL constants (whitespace-only, parity-verified) --- gateway/hosted_room_driver.py | 16 +- gateway/hosted_room_policy_checkpoint.py | 42 ++--- gateway/hosted_room_replicas.py | 44 ++--- gateway/hosted_rooms.py | 216 ++++++++--------------- 4 files changed, 119 insertions(+), 199 deletions(-) diff --git a/gateway/hosted_room_driver.py b/gateway/hosted_room_driver.py index 87723a8406..256867f34b 100644 --- a/gateway/hosted_room_driver.py +++ b/gateway/hosted_room_driver.py @@ -592,10 +592,9 @@ def acquire_lease( row = conn.execute(_SELECT_LEASE, (room_id,)).fetchone() if row is None: conn.execute( - """INSERT INTO hosted_room_driver_leases ( - room_id, gateway_id, authority_epoch, process_generation, lease_generation, - expires_at, acquired_at, updated_at, released_at - ) VALUES (?, ?, ?, ?, 1, ?, ?, ?, NULL)""", + """INSERT INTO hosted_room_driver_leases ( room_id, gateway_id, authority_epoch, + process_generation, lease_generation, expires_at, acquired_at, updated_at, + released_at ) VALUES (?, ?, ?, ?, 1, ?, ?, ?, NULL)""", (room_id, gateway_id, authority_epoch, process_generation, expires_at, now, now), ) return _lease_from_row(conn.execute(_SELECT_LEASE, (room_id,)).fetchone()) @@ -1044,11 +1043,10 @@ def prune_published_terminal_tasks( if publications is None: return 0 rows = conn.execute( - """SELECT t.task_id, t.terminal_at FROM hosted_room_driver_tasks t - WHERE t.room_id=? AND t.status IN ('settled', 'failed', 'cancelled') - AND EXISTS (SELECT 1 FROM hosted_room_policy_publications p - WHERE p.room_id=t.room_id AND p.task_id=t.task_id - AND p.kind IN ('turn.settled', 'turn.failed', 'turn.cancelled')) + """SELECT t.task_id, t.terminal_at FROM hosted_room_driver_tasks t WHERE t.room_id=? + AND t.status IN ('settled', 'failed', 'cancelled') AND EXISTS (SELECT 1 FROM + hosted_room_policy_publications p WHERE p.room_id=t.room_id AND p.task_id=t.task_id + AND p.kind IN ('turn.settled', 'turn.failed', 'turn.cancelled')) ORDER BY t.terminal_at DESC, t.task_id ASC""", (room_id,), ).fetchall() diff --git a/gateway/hosted_room_policy_checkpoint.py b/gateway/hosted_room_policy_checkpoint.py index 590f109ed8..0d55c694a9 100644 --- a/gateway/hosted_room_policy_checkpoint.py +++ b/gateway/hosted_room_policy_checkpoint.py @@ -39,10 +39,9 @@ _SCHEMA_DDL = ( """CREATE TABLE IF NOT EXISTS hosted_room_policy_watermarks ( room_id TEXT NOT NULL, thread_id TEXT NOT NULL, member_id TEXT NOT NULL, seen_through_seq INTEGER NOT NULL, PRIMARY KEY(room_id, thread_id, member_id))""", - """CREATE TABLE IF NOT EXISTS hosted_room_policy_publications ( - room_id TEXT NOT NULL, task_id TEXT NOT NULL, kind TEXT NOT NULL, - execution_generation INTEGER NOT NULL DEFAULT 0, seq INTEGER NOT NULL, - PRIMARY KEY(room_id, task_id, kind, execution_generation))""", + """CREATE TABLE IF NOT EXISTS hosted_room_policy_publications ( room_id TEXT NOT NULL, + task_id TEXT NOT NULL, kind TEXT NOT NULL, execution_generation INTEGER NOT NULL DEFAULT 0, + seq INTEGER NOT NULL, PRIMARY KEY(room_id, task_id, kind, execution_generation))""", # Transcript stores only references into the already bounded room log, so # prompt payloads are never duplicated outside room byte limits. """CREATE TABLE IF NOT EXISTS hosted_room_policy_transcript ( @@ -116,9 +115,8 @@ class HostedRoomPolicyCheckpoint: conn: sqlite3.Connection, *, event: Mapping[str, Any], thread_id: str, discussion_event_id: str ) -> None: conn.execute( - """INSERT OR IGNORE INTO hosted_room_policy_events( - room_id, thread_id, discussion_event_id, seq, event_json - ) VALUES (?, ?, ?, ?, ?)""", + """INSERT OR IGNORE INTO hosted_room_policy_events( room_id, thread_id, + discussion_event_id, seq, event_json ) VALUES (?, ?, ?, ?, ?)""", ( event["room_id"], thread_id, discussion_event_id, int(event["seq"]), json.dumps(dict(event), ensure_ascii=True, sort_keys=True, separators=(",", ":")), @@ -131,11 +129,10 @@ class HostedRoomPolicyCheckpoint: settled_seq: int | None = None, ) -> None: conn.execute( - """INSERT INTO hosted_room_policy_transcript( - room_id, thread_id, seq, kind, settled_seq - ) VALUES (?, ?, ?, ?, ?) - ON CONFLICT(room_id, thread_id, seq) DO UPDATE SET - settled_seq=COALESCE(excluded.settled_seq, hosted_room_policy_transcript.settled_seq)""", + """INSERT INTO hosted_room_policy_transcript( room_id, thread_id, seq, kind, settled_seq + ) VALUES (?, ?, ?, ?, ?) ON CONFLICT(room_id, thread_id, seq) DO UPDATE SET + settled_seq=COALESCE(excluded.settled_seq, + hosted_room_policy_transcript.settled_seq)""", (event["room_id"], thread_id, int(event["seq"]), str(event["kind"]), settled_seq), ) if event["kind"] in {"message.user", "message.member"}: @@ -216,12 +213,11 @@ class HostedRoomPolicyCheckpoint: if not thread_id or not event_id: return conn.execute( - """INSERT INTO hosted_room_policy_threads( - room_id, thread_id, discussion_event_id, latest_user_seq, completed - ) VALUES (?, ?, ?, ?, 0) - ON CONFLICT(room_id, thread_id) DO UPDATE SET - discussion_event_id=excluded.discussion_event_id, - latest_user_seq=excluded.latest_user_seq, completed=0""", + """INSERT INTO hosted_room_policy_threads( room_id, thread_id, discussion_event_id, + latest_user_seq, completed ) VALUES (?, ?, ?, ?, 0) + ON CONFLICT(room_id, thread_id) DO UPDATE + SET discussion_event_id=excluded.discussion_event_id, + latest_user_seq=excluded.latest_user_seq, completed=0""", (room_id, thread_id, event_id, int(event["seq"])), ) self._store_active_event(conn, event=event, thread_id=thread_id, discussion_event_id=event_id) @@ -253,9 +249,8 @@ class HostedRoomPolicyCheckpoint: ) if task_id: conn.execute( - """INSERT OR IGNORE INTO hosted_room_policy_publications( - room_id, task_id, kind, execution_generation, seq - ) VALUES (?, ?, ?, ?, ?)""", + """INSERT OR IGNORE INTO hosted_room_policy_publications( room_id, task_id, kind, + execution_generation, seq ) VALUES (?, ?, ?, ?, ?)""", (room_id, task_id, kind, execution_generation, seq), ) member_id = str(payload.get("member_id") or "") @@ -327,9 +322,8 @@ class HostedRoomPolicyCheckpoint: """Create the room cursor if absent, backfill the transcript once, return through_seq.""" _require_room(conn, room_id) conn.execute( - """INSERT OR IGNORE INTO hosted_room_policy_cursors( - room_id, through_seq, stopped_through_seq, updated_at - ) VALUES (?, 0, 0, 0)""", + """INSERT OR IGNORE INTO hosted_room_policy_cursors( room_id, through_seq, + stopped_through_seq, updated_at ) VALUES (?, 0, 0, 0)""", (room_id,), ) row = conn.execute( diff --git a/gateway/hosted_room_replicas.py b/gateway/hosted_room_replicas.py index 4944293b3d..45c7cab3de 100644 --- a/gateway/hosted_room_replicas.py +++ b/gateway/hosted_room_replicas.py @@ -60,22 +60,18 @@ class ReplicaEpochRegressionError(ReplicaError): def _initialize_replica_schema(conn: sqlite3.Connection) -> None: conn.execute( - """CREATE TABLE IF NOT EXISTS hosted_room_replicas ( - room_id TEXT PRIMARY KEY, name TEXT NOT NULL, members_json TEXT NOT NULL, - authority_gateway_id TEXT NOT NULL, + """CREATE TABLE IF NOT EXISTS hosted_room_replicas ( room_id TEXT PRIMARY KEY, + name TEXT NOT NULL, members_json TEXT NOT NULL, authority_gateway_id TEXT NOT NULL, authority_epoch INTEGER NOT NULL CHECK (authority_epoch >= 1), last_seq INTEGER NOT NULL DEFAULT 0 CHECK (last_seq >= 0), latest_seq INTEGER NOT NULL DEFAULT 0, event_bytes INTEGER NOT NULL DEFAULT 0, - created_at REAL NOT NULL, updated_at REAL NOT NULL - )""" + created_at REAL NOT NULL, updated_at REAL NOT NULL )""" ) conn.execute( - """CREATE TABLE IF NOT EXISTS hosted_room_replica_events ( - room_id TEXT NOT NULL, seq INTEGER NOT NULL CHECK (seq >= 1), event_id TEXT NOT NULL, - kind TEXT NOT NULL, actor_json TEXT NOT NULL, authority_epoch INTEGER, - payload_json TEXT NOT NULL, created_at REAL NOT NULL, - PRIMARY KEY (room_id, seq) - )""" + """CREATE TABLE IF NOT EXISTS hosted_room_replica_events ( room_id TEXT NOT NULL, + seq INTEGER NOT NULL CHECK (seq >= 1), event_id TEXT NOT NULL, kind TEXT NOT NULL, + actor_json TEXT NOT NULL, authority_epoch INTEGER, payload_json TEXT NOT NULL, + created_at REAL NOT NULL, PRIMARY KEY (room_id, seq) )""" ) @@ -208,10 +204,9 @@ def ingest_page( latest_seq = new_last if row is None: conn.execute( - """INSERT INTO hosted_room_replicas - (room_id, name, members_json, authority_gateway_id, authority_epoch, last_seq, latest_seq, - event_bytes, created_at, updated_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + """INSERT INTO hosted_room_replicas (room_id, name, members_json, + authority_gateway_id, authority_epoch, last_seq, latest_seq, event_bytes, + created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", ( room_id, room_name, members_json, authority["gateway_id"], authority["epoch"], new_last, max(latest_seq, new_last), added_bytes, now, now, @@ -219,10 +214,9 @@ def ingest_page( ) else: conn.execute( - """UPDATE hosted_room_replicas - SET name=?, members_json=?, authority_gateway_id=?, authority_epoch=?, last_seq=?, - latest_seq=?, event_bytes=event_bytes+?, updated_at=? - WHERE room_id=?""", + """UPDATE hosted_room_replicas SET name=?, members_json=?, authority_gateway_id=?, + authority_epoch=?, last_seq=?, latest_seq=?, event_bytes=event_bytes+?, + updated_at=? WHERE room_id=?""", ( room_name, members_json, authority["gateway_id"], authority["epoch"], new_last, max(latest_seq, new_last), added_bytes, now, room_id, @@ -299,10 +293,9 @@ def promote_replica( claim_bytes = utf8_len(claim_event_id, "authority.claimed", claim_actor_json, claim_payload_json) conn.execute( - """INSERT INTO hosted_rooms - (room_id, name, members_json, authority_gateway_id, authority_epoch, next_seq, event_bytes, - revision, created_at, updated_at, disbanded_at) - VALUES (?, ?, ?, ?, ?, ?, ?, 1, ?, ?, NULL)""", + """INSERT INTO hosted_rooms (room_id, name, members_json, authority_gateway_id, + authority_epoch, next_seq, event_bytes, revision, created_at, updated_at, + disbanded_at) VALUES (?, ?, ?, ?, ?, ?, ?, 1, ?, ?, NULL)""", ( room_id, replica["name"], replica["members_json"], local_gateway, target_epoch, claim_seq + 1, int(replica["event_bytes"]) + claim_bytes, now, now, @@ -386,9 +379,8 @@ def demote_room( ), ) conn.execute( - """UPDATE hosted_rooms - SET authority_gateway_id=?, authority_epoch=?, next_seq=next_seq+1, revision=revision+1, updated_at=? - WHERE room_id=?""", + """UPDATE hosted_rooms SET authority_gateway_id=?, authority_epoch=?, + next_seq=next_seq+1, revision=revision+1, updated_at=? WHERE room_id=?""", (observed_gateway_id, observed_epoch, now, room_id), ) return { diff --git a/gateway/hosted_rooms.py b/gateway/hosted_rooms.py index 618ebf440d..b5b5599ee3 100644 --- a/gateway/hosted_rooms.py +++ b/gateway/hosted_rooms.py @@ -102,68 +102,35 @@ _REMOTE_RUNS_BODY = """ """ # Executed in this exact order on first open / migration. _SCHEMA_DDL = ( - """CREATE TABLE IF NOT EXISTS hosted_rooms ( - room_id TEXT PRIMARY KEY, - name TEXT NOT NULL, - members_json TEXT NOT NULL, - authority_gateway_id TEXT NOT NULL, - authority_epoch INTEGER NOT NULL DEFAULT 1 CHECK (authority_epoch >= 1), - next_seq INTEGER NOT NULL DEFAULT 1 CHECK (next_seq >= 1), - event_bytes INTEGER NOT NULL DEFAULT 0 CHECK (event_bytes >= 0), - revision INTEGER NOT NULL DEFAULT 1 CHECK (revision >= 1), - created_at REAL NOT NULL, - updated_at REAL NOT NULL, - disbanded_at REAL - )""", - """CREATE TABLE IF NOT EXISTS hosted_room_events ( - room_id TEXT NOT NULL, - seq INTEGER NOT NULL CHECK (seq >= 1), - event_id TEXT NOT NULL, - kind TEXT NOT NULL, - actor_json TEXT NOT NULL, - authority_epoch INTEGER CHECK (authority_epoch IS NULL OR authority_epoch >= 1), - payload_json TEXT NOT NULL, - created_at REAL NOT NULL, - PRIMARY KEY (room_id, seq), - UNIQUE (room_id, event_id), - FOREIGN KEY (room_id) REFERENCES hosted_rooms(room_id) - )""", - """CREATE TABLE IF NOT EXISTS hosted_room_retired_ids ( - room_id TEXT PRIMARY KEY, - retired_at REAL NOT NULL - )""", - """CREATE TABLE IF NOT EXISTS hosted_room_links ( - room_id TEXT NOT NULL, - member_id TEXT NOT NULL, - target_url TEXT NOT NULL, - target_profile TEXT NOT NULL, - grant TEXT NOT NULL, - catalog_json TEXT NOT NULL, - cancellation_scope_id TEXT NOT NULL, - trace_id TEXT NOT NULL, - transport_security TEXT NOT NULL, - status TEXT NOT NULL DEFAULT 'ready', - updated_at REAL NOT NULL, - PRIMARY KEY (room_id, member_id) - )""", + """CREATE TABLE IF NOT EXISTS hosted_rooms ( room_id TEXT PRIMARY KEY, name TEXT NOT NULL, + members_json TEXT NOT NULL, authority_gateway_id TEXT NOT NULL, + authority_epoch INTEGER NOT NULL DEFAULT 1 CHECK (authority_epoch >= 1), + next_seq INTEGER NOT NULL DEFAULT 1 CHECK (next_seq >= 1), + event_bytes INTEGER NOT NULL DEFAULT 0 CHECK (event_bytes >= 0), + revision INTEGER NOT NULL DEFAULT 1 CHECK (revision >= 1), created_at REAL NOT NULL, + updated_at REAL NOT NULL, disbanded_at REAL )""", + """CREATE TABLE IF NOT EXISTS hosted_room_events ( room_id TEXT NOT NULL, + seq INTEGER NOT NULL CHECK (seq >= 1), event_id TEXT NOT NULL, kind TEXT NOT NULL, + actor_json TEXT NOT NULL, + authority_epoch INTEGER CHECK (authority_epoch IS NULL OR authority_epoch >= 1), + payload_json TEXT NOT NULL, created_at REAL NOT NULL, PRIMARY KEY (room_id, seq), + UNIQUE (room_id, event_id), FOREIGN KEY (room_id) REFERENCES hosted_rooms(room_id) )""", + """CREATE TABLE IF NOT EXISTS hosted_room_retired_ids ( room_id TEXT PRIMARY KEY, + retired_at REAL NOT NULL )""", + """CREATE TABLE IF NOT EXISTS hosted_room_links ( room_id TEXT NOT NULL, + member_id TEXT NOT NULL, target_url TEXT NOT NULL, target_profile TEXT NOT NULL, + grant TEXT NOT NULL, catalog_json TEXT NOT NULL, cancellation_scope_id TEXT NOT NULL, + trace_id TEXT NOT NULL, transport_security TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'ready', updated_at REAL NOT NULL, + PRIMARY KEY (room_id, member_id) )""", f"CREATE TABLE IF NOT EXISTS hosted_room_remote_runs ({_REMOTE_RUNS_BODY})", - """CREATE TABLE IF NOT EXISTS hosted_room_revoked_grants ( - scope_key TEXT PRIMARY KEY, - expires_at REAL NOT NULL, - revoked_before REAL NOT NULL - )""", - """CREATE TABLE IF NOT EXISTS hosted_room_peer_reservations ( - room_id TEXT NOT NULL, - member_id TEXT NOT NULL, - target_profile TEXT NOT NULL, - authority_gateway_id TEXT NOT NULL, - authority_epoch INTEGER NOT NULL CHECK (authority_epoch >= 1), - expires_at REAL NOT NULL, - revoked_at REAL, - created_at REAL NOT NULL, - updated_at REAL NOT NULL, - PRIMARY KEY (room_id, member_id, target_profile) - )""", + """CREATE TABLE IF NOT EXISTS hosted_room_revoked_grants ( scope_key TEXT PRIMARY KEY, + expires_at REAL NOT NULL, revoked_before REAL NOT NULL )""", + """CREATE TABLE IF NOT EXISTS hosted_room_peer_reservations ( room_id TEXT NOT NULL, + member_id TEXT NOT NULL, target_profile TEXT NOT NULL, authority_gateway_id TEXT NOT NULL, + authority_epoch INTEGER NOT NULL CHECK (authority_epoch >= 1), expires_at REAL NOT NULL, + revoked_at REAL, created_at REAL NOT NULL, updated_at REAL NOT NULL, + PRIMARY KEY (room_id, member_id, target_profile) )""", ) # (table, required columns) in the order _schema_is_current probes them. _REQUIRED_COLUMNS = ( @@ -476,17 +443,10 @@ def _migrate_legacy_columns(conn: sqlite3.Connection) -> None: conn.execute("ALTER TABLE hosted_room_events ADD COLUMN authority_epoch INTEGER") if backfill_event_bytes: conn.execute( - """UPDATE hosted_rooms - SET event_bytes=COALESCE(( - SELECT SUM( - length(CAST(event_id AS BLOB)) + - length(CAST(kind AS BLOB)) + - length(CAST(actor_json AS BLOB)) + - length(CAST(payload_json AS BLOB)) - ) - FROM hosted_room_events - WHERE hosted_room_events.room_id=hosted_rooms.room_id - ), 0)""" + """UPDATE hosted_rooms SET event_bytes=COALESCE(( SELECT SUM( length(CAST(event_id AS + BLOB)) + length(CAST(kind AS BLOB)) + length(CAST(actor_json AS BLOB)) + + length(CAST(payload_json AS BLOB)) ) FROM hosted_room_events WHERE + hosted_room_events.room_id=hosted_rooms.room_id ), 0)""" ) @@ -775,11 +735,9 @@ def list_room_link_records(db_path: Path | str) -> list[dict[str, Any]]: """Return private RoomLink records without logging or formatting grants.""" with _transaction(db_path) as conn: rows = conn.execute( - """SELECT room_id, member_id, target_url, target_profile, grant, - catalog_json, cancellation_scope_id, trace_id, - transport_security, status, updated_at - FROM hosted_room_links - ORDER BY room_id, member_id""" + """SELECT room_id, member_id, target_url, target_profile, grant, catalog_json, + cancellation_scope_id, trace_id, transport_security, status, + updated_at FROM hosted_room_links ORDER BY room_id, member_id""" ).fetchall() return [dict(row) for row in rows] @@ -798,21 +756,15 @@ def upsert_room_link_record( if count >= max_links: raise HostedRoomError("too many stored room links") conn.execute( - """INSERT INTO hosted_room_links( - room_id, member_id, target_url, target_profile, grant, - catalog_json, cancellation_scope_id, trace_id, - transport_security, status, updated_at - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(room_id, member_id) DO UPDATE SET - target_url=excluded.target_url, - target_profile=excluded.target_profile, - grant=excluded.grant, - catalog_json=excluded.catalog_json, - cancellation_scope_id=excluded.cancellation_scope_id, - trace_id=excluded.trace_id, - transport_security=excluded.transport_security, - status=excluded.status, - updated_at=excluded.updated_at""", + """INSERT INTO hosted_room_links( room_id, member_id, target_url, target_profile, grant, + catalog_json, cancellation_scope_id, trace_id, transport_security, status, + updated_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(room_id, member_id) DO UPDATE SET target_url=excluded.target_url, + target_profile=excluded.target_profile, grant=excluded.grant, + catalog_json=excluded.catalog_json, + cancellation_scope_id=excluded.cancellation_scope_id, trace_id=excluded.trace_id, + transport_security=excluded.transport_security, status=excluded.status, + updated_at=excluded.updated_at""", ( record["room_id"], record["member_id"], record["target_url"], record["target_profile"], record["grant"], record["catalog_json"], @@ -868,14 +820,11 @@ def revoke_room_grant_scope( with _transaction(db_path, immediate=True) as conn: conn.execute("DELETE FROM hosted_room_revoked_grants WHERE expires_at<=?", (timestamp,)) conn.execute( - """INSERT INTO hosted_room_revoked_grants( - scope_key, expires_at, revoked_before - ) VALUES (?, ?, ?) - ON CONFLICT(scope_key) DO UPDATE SET - expires_at=MAX(hosted_room_revoked_grants.expires_at, - excluded.expires_at), - revoked_before=MAX(hosted_room_revoked_grants.revoked_before, - excluded.revoked_before)""", + """INSERT INTO hosted_room_revoked_grants( scope_key, expires_at, revoked_before ) + VALUES (?, ?, ?) ON CONFLICT(scope_key) DO UPDATE + SET expires_at=MAX(hosted_room_revoked_grants.expires_at, excluded.expires_at), + revoked_before=MAX(hosted_room_revoked_grants.revoked_before, + excluded.revoked_before)""", (scope_key, expiry, timestamp), ) conn.execute( @@ -946,17 +895,14 @@ def reserve_peer_room( if existing is not None and _reservation_superseded(existing, gateway_id, epoch): raise AuthorityConflictError("peer room reservation authority changed") conn.execute( - """INSERT INTO hosted_room_peer_reservations( - room_id, member_id, target_profile, authority_gateway_id, - authority_epoch, expires_at, revoked_at, created_at, updated_at - ) VALUES (?, ?, ?, ?, ?, ?, NULL, ?, ?) - ON CONFLICT(room_id, member_id, target_profile) DO UPDATE SET - authority_gateway_id=excluded.authority_gateway_id, - authority_epoch=excluded.authority_epoch, - expires_at=MAX(hosted_room_peer_reservations.expires_at, - excluded.expires_at), - revoked_at=NULL, - updated_at=excluded.updated_at""", + """INSERT INTO hosted_room_peer_reservations( room_id, member_id, target_profile, + authority_gateway_id, authority_epoch, expires_at, revoked_at, created_at, + updated_at ) VALUES (?, ?, ?, ?, ?, ?, NULL, ?, ?) + ON CONFLICT(room_id, member_id, target_profile) DO UPDATE + SET authority_gateway_id=excluded.authority_gateway_id, + authority_epoch=excluded.authority_epoch, + expires_at=MAX(hosted_room_peer_reservations.expires_at, excluded.expires_at), + revoked_at=NULL, updated_at=excluded.updated_at""", (*values, expiry, timestamp, timestamp), ) @@ -1031,12 +977,10 @@ def upsert_remote_run_receipt( ) return conn.execute( - """INSERT INTO hosted_room_remote_runs( - room_id, home_install_id, authority_gateway_id, - authority_epoch, member_id, target_install_id, - target_profile, task_id, execution_generation, run_id, - session_id, created_at, updated_at - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + """INSERT INTO hosted_room_remote_runs( room_id, home_install_id, authority_gateway_id, + authority_epoch, member_id, target_install_id, target_profile, task_id, + execution_generation, run_id, session_id, created_at, updated_at ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", (*immutable, timestamp, timestamp), ) @@ -1097,9 +1041,8 @@ def _adopt_legacy_room( ), ) adopted = conn.execute( - """UPDATE hosted_rooms - SET members_json=?, authority_gateway_id=?, authority_epoch=?, - next_seq=next_seq+1, revision=revision+1, event_bytes=event_bytes+?, updated_at=? + """UPDATE hosted_rooms SET members_json=?, authority_gateway_id=?, authority_epoch=?, + next_seq=next_seq+1, revision=revision+1, event_bytes=event_bytes+?, updated_at=? WHERE room_id=? AND authority_gateway_id='legacy' AND authority_epoch=? AND next_seq=? AND disbanded_at IS NULL""", ( @@ -1174,9 +1117,8 @@ def create_room( (room_id, name, members_json, authority_gateway_id, now, now), ) row = conn.execute( - """SELECT room_id, name, members_json, authority_gateway_id, - authority_epoch, revision, created_at, updated_at - FROM hosted_rooms WHERE room_id=?""", + """SELECT room_id, name, members_json, authority_gateway_id, authority_epoch, revision, + created_at, updated_at FROM hosted_rooms WHERE room_id=?""", (room_id,), ).fetchone() if row is None: # pragma: no cover - guarded by the insert above @@ -1235,9 +1177,8 @@ def rename_room( event_bytes = _prepare_event(conn, room, event_id, "room.renamed", actor_json, payload_json) # Rename updates the room row before inserting its event (order is load-bearing). conn.execute( - """UPDATE hosted_rooms - SET name=?, next_seq=?, event_bytes=event_bytes+?, revision=revision+1, updated_at=? - WHERE room_id=?""", + """UPDATE hosted_rooms SET name=?, next_seq=?, event_bytes=event_bytes+?, + revision=revision+1, updated_at=? WHERE room_id=?""", (name, seq + 1, event_bytes, now, room_id), ) conn.execute( @@ -1471,12 +1412,10 @@ def claim_authority( ), ) updated = conn.execute( - """UPDATE hosted_rooms - SET authority_gateway_id=?, authority_epoch=authority_epoch+1, - next_seq=next_seq+1, event_bytes=event_bytes+?, - revision=revision+1, updated_at=? - WHERE room_id=? AND disbanded_at IS NULL AND authority_gateway_id=? - AND authority_epoch=?""", + """UPDATE hosted_rooms SET authority_gateway_id=?, + authority_epoch=authority_epoch+1, next_seq=next_seq+1, + event_bytes=event_bytes+?, revision=revision+1, updated_at=? WHERE room_id=? + AND disbanded_at IS NULL AND authority_gateway_id=? AND authority_epoch=?""", (new_gateway_id, claim_bytes, now, room_id, expected_gateway_id, expected_epoch), ) if updated.rowcount != 1: @@ -1484,9 +1423,8 @@ def claim_authority( idempotent = False existing_event = _load_event(conn, room_id, event_id) state_row = conn.execute( - """SELECT room_id, name, members_json, authority_gateway_id, - authority_epoch, next_seq, revision, created_at, updated_at - FROM hosted_rooms WHERE room_id=?""", + """SELECT room_id, name, members_json, authority_gateway_id, authority_epoch, next_seq, + revision, created_at, updated_at FROM hosted_rooms WHERE room_id=?""", (room_id,), ).fetchone() if state_row is None: # pragma: no cover - room exists in this transaction @@ -1554,11 +1492,9 @@ def disband_room( ), ) updated = conn.execute( - """UPDATE hosted_rooms - SET disbanded_at=?, updated_at=?, revision=revision+1, - next_seq=next_seq+1, event_bytes=event_bytes+? - WHERE room_id=? AND disbanded_at IS NULL AND authority_gateway_id=? - AND authority_epoch=?""", + """UPDATE hosted_rooms SET disbanded_at=?, updated_at=?, revision=revision+1, + next_seq=next_seq+1, event_bytes=event_bytes+? WHERE room_id=? + AND disbanded_at IS NULL AND authority_gateway_id=? AND authority_epoch=?""", (now, now, disband_bytes, room_id, expected_gateway_id, expected_epoch), ) if updated.rowcount != 1: