Revert "refactor(gateway/hosted_rooms): reflow multi-line SQL constants (whitespace-only, parity-verified)"
This reverts commit fa827f596e.
This commit is contained in:
@@ -630,9 +630,10 @@ 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())
|
||||
@@ -1041,10 +1042,11 @@ 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,9 +39,10 @@ _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 (
|
||||
@@ -129,8 +130,9 @@ 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=(",", ":")),
|
||||
@@ -143,10 +145,11 @@ 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"}:
|
||||
@@ -227,11 +230,12 @@ 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)
|
||||
@@ -263,8 +267,9 @@ 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 "")
|
||||
@@ -324,8 +329,9 @@ 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,18 +60,22 @@ 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)
|
||||
)"""
|
||||
)
|
||||
|
||||
|
||||
@@ -301,9 +305,10 @@ 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,
|
||||
@@ -387,8 +392,9 @@ 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 {
|
||||
|
||||
+128
-68
@@ -102,35 +102,68 @@ _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 = (
|
||||
@@ -464,10 +497,17 @@ def _migrate_legacy_columns(conn: sqlite3.Connection) -> None:
|
||||
conn.execute(ddl)
|
||||
if "event_bytes" not in columns["hosted_rooms"]:
|
||||
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)"""
|
||||
)
|
||||
|
||||
|
||||
@@ -748,9 +788,11 @@ 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]
|
||||
|
||||
@@ -769,15 +811,21 @@ 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"],
|
||||
@@ -833,11 +881,14 @@ 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(
|
||||
@@ -908,14 +959,17 @@ 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),
|
||||
)
|
||||
|
||||
@@ -991,10 +1045,12 @@ 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),
|
||||
)
|
||||
|
||||
@@ -1055,8 +1111,9 @@ 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""",
|
||||
(
|
||||
@@ -1188,8 +1245,9 @@ 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(
|
||||
@@ -1520,9 +1578,11 @@ 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