refactor(gateway/hosted_rooms): reflow multi-line SQL constants (whitespace-only, parity-verified)
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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 {
|
||||
|
||||
+76
-140
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user