Merge branch 'simp/r2-gw-rest' into simp/integration2

This commit is contained in:
Teknium
2026-09-02 17:27:35 -07:00
6 changed files with 551 additions and 451 deletions

View File

@@ -432,13 +432,13 @@ def _validate_member_message(
raise DiscussionValidationError("message.member text must be a non-pass string")
member = _member_by_id(room, payload.get("member_id"))
peer = _peer_target(member)
expected_connection = peer.get("peer_id") if peer else None
if (
actor.get("kind") != "member"
or actor.get("id") != member.member_id
or actor.get("profile") != member.profile
or actor.get("connection_id") != expected_connection
):
expected = {
"kind": "member",
"id": member.member_id,
"profile": member.profile,
"connection_id": peer.get("peer_id") if peer else None,
}
if any(actor.get(key) != value for key, value in expected.items()):
raise DiscussionValidationError("message.member actor does not match roster")
return payload
@@ -510,15 +510,14 @@ def _validate_event(raw: Any, *, room: DiscussionRoom, previous_seq: int) -> _Va
if seq <= previous_seq:
raise DiscussionValidationError("room events must be in strict sequence order")
event_id = _identifier(raw.get("event_id"), label="event_id")
kind = raw.get("kind")
if not isinstance(kind, str):
raise DiscussionValidationError("event kind must be a string")
actor = raw.get("actor")
if not isinstance(actor, Mapping):
raise DiscussionValidationError("event actor must be an object")
payload = raw.get("payload")
if not isinstance(payload, Mapping):
raise DiscussionValidationError("event payload must be an object")
kind, actor, payload = raw.get("kind"), raw.get("actor"), raw.get("payload")
for value, expected, message in (
(kind, str, "event kind must be a string"),
(actor, Mapping, "event actor must be an object"),
(payload, Mapping, "event payload must be an object"),
):
if not isinstance(value, expected):
raise DiscussionValidationError(message)
if kind in _EPOCH_STAMPED_KINDS and raw.get("authority_epoch") != room.authority_epoch:
raise DiscussionValidationError(f"{kind} authority epoch does not match the room")
validator = _EVENT_VALIDATORS.get(kind)
@@ -576,11 +575,9 @@ def _derive_member_watermarks(events: Sequence[_ValidatedEvent]) -> dict[tuple[s
watermark = int(event.payload["seen_through_seq"])
if event.kind == "turn.settled" and not event.payload["passed"]:
message = messages_by_id.get(str(event.payload["message_event_id"]))
if (
message is None
or message.payload.get("task_id") != task_id
or message.payload.get("member_id") != event.payload.get("member_id")
or message.payload.get("thread_id") != event.payload.get("thread_id")
if message is None or any(
message.payload.get(field) != event.payload.get(field)
for field in ("task_id", "member_id", "thread_id")
):
raise DiscussionValidationError("turn.settled references no matching member message")
watermark = max(watermark, message.seq)
@@ -1021,12 +1018,9 @@ def plan_publication(
raise DiscussionValidationError("task member is not in the frozen roster")
if status not in _TERMINAL_EFFECTS:
raise DiscussionValidationError("invalid terminal publication status")
if status == "deferred" and (
isinstance(execution_generation, bool)
or not isinstance(execution_generation, int)
or execution_generation < 1
):
raise DiscussionValidationError("deferred publication requires an execution generation")
if status == "deferred":
message = "deferred publication requires an execution generation"
common.positive_int(execution_generation, error=DiscussionValidationError, message=message)
newer_same_thread = any(
event.kind == "message.user"

View File

@@ -69,27 +69,13 @@ def _task_update(set_clause: str, fence: str) -> str:
return f"UPDATE hosted_room_driver_tasks SET {set_clause} WHERE room_id=? AND task_id=? AND {fence}"
def _settle_sql(status: str) -> str:
"""Terminal settlement of a task fenced on status + both generations."""
return _task_update(_SETTLE_SET, f"status='{status}' AND {_GENERATION_FENCE}")
def _generation_update(set_clause: str, status: str) -> str:
"""Transition fenced on ``status`` + both generations (terminal settlements and the recovery family)."""
return _task_update(set_clause, f"status='{status}' AND {_GENERATION_FENCE}")
_SETTLE_RUNNING_SQL = _settle_sql("running") + f" AND {_RUN_FENCE}"
_SETTLE_STOPPING_SQL = _settle_sql("stopping")
_SETTLE_INDETERMINATE_SQL = _settle_sql("indeterminate")
_CANCEL_INDETERMINATE_SQL = _task_update(_CANCEL_SET, f"status='indeterminate' AND {_GENERATION_FENCE}")
_REQUEUE_INDETERMINATE_SQL = _task_update(
f"{_REQUEUE_SET}, started_at=NULL, indeterminate_at=NULL, updated_at=?",
f"status='indeterminate' AND {_GENERATION_FENCE}",
)
_DEFER_SQL = _task_update(
"status='deferred', result_json=?, terminal_at=?, updated_at=?",
f"status='indeterminate' AND {_GENERATION_FENCE}",
)
_REQUEUE_DEFERRED_SQL = _task_update(
f"{_REQUEUE_SET}, result_json=NULL, started_at=NULL, terminal_at=NULL, indeterminate_at=NULL, updated_at=?",
f"status='deferred' AND {_GENERATION_FENCE}",
)
_SETTLE_RUNNING_SQL = _generation_update(_SETTLE_SET, "running") + f" AND {_RUN_FENCE}"
_SETTLE_STOPPING_SQL = _generation_update(_SETTLE_SET, "stopping")
_REQUEUE_RUNNING_SQL = _task_update(
f"{_REQUEUE_SET}, started_at=NULL, updated_at=?",
f"status='running' AND {_GENERATION_FENCE} AND {_RUN_FENCE}",
@@ -104,6 +90,32 @@ _COMPLETE_STOP_SQL = _task_update(
"status='stopping' AND cancel_id=? AND cancel_generation=?",
)
# Lease-first recovery transitions: name -> (fenced status, SET clause, generation-guard stale message,
# row stale message); the UPDATE is _generation_update(set_clause, status).
_INDETERMINATE_STALE = "indeterminate task generation changed"
_GENERATION_TRANSITIONS = {
"resolve": (
"indeterminate", _SETTLE_SET, _INDETERMINATE_STALE, "indeterminate task changed during reconciliation",
),
"resolve_cancel": (
"indeterminate", _CANCEL_SET, "indeterminate cancellation proof is stale",
"indeterminate cancellation proof lost its fence",
),
"requeue": (
"indeterminate", f"{_REQUEUE_SET}, started_at=NULL, indeterminate_at=NULL, updated_at=?",
_INDETERMINATE_STALE, "indeterminate task changed during requeue",
),
"defer": (
"indeterminate", "status='deferred', result_json=?, terminal_at=?, updated_at=?",
_INDETERMINATE_STALE, "indeterminate task changed during deferral",
),
"requeue_deferred": (
"deferred",
f"{_REQUEUE_SET}, result_json=NULL, started_at=NULL, terminal_at=NULL, indeterminate_at=NULL, updated_at=?",
"deferred task generation changed", "deferred task changed during requeue",
),
}
class DriverStateError(ValueError):
"""Base class for invalid or conflicting driver-state operations."""
@@ -141,26 +153,24 @@ def _identifier(value: Any, *, label: str) -> str:
return identifier(value, label=label, error=DriverValidationError, max_chars=MAX_IDENTIFIER_CHARS)
def _timestamp(clock: Clock) -> float:
if not callable(clock):
raise DriverValidationError("clock must be callable")
def _finite(compute: Callable[[], Any], message: str, *, positive: bool = False) -> float:
try:
value = float(clock())
value = float(compute())
except (TypeError, ValueError, OverflowError) as exc:
raise DriverValidationError("clock must return a finite number") from exc
if not math.isfinite(value):
raise DriverValidationError("clock must return a finite number")
raise DriverValidationError(message) from exc
if not math.isfinite(value) or (positive and value <= 0):
raise DriverValidationError(message)
return value
def _timestamp(clock: Clock) -> float:
if not callable(clock):
raise DriverValidationError("clock must be callable")
return _finite(clock, "clock must return a finite number")
def _ttl(value: Any) -> float:
try:
ttl = float(value)
except (TypeError, ValueError, OverflowError) as exc:
raise DriverValidationError("ttl_seconds must be a finite positive number") from exc
if not math.isfinite(ttl) or ttl <= 0:
raise DriverValidationError("ttl_seconds must be a finite positive number")
return ttl
return _finite(lambda: value, "ttl_seconds must be a finite positive number", positive=True)
def _expiry(now: float, ttl: float) -> float:
@@ -515,14 +525,6 @@ def _generations_match(row: sqlite3.Row, status: str, execution_generation: int,
)
def _generation_guard(status: str, execution_generation: int, cancel_generation: int, stale: str):
def guard(row: sqlite3.Row) -> None:
if not _generations_match(row, status, execution_generation, cancel_generation):
raise StaleTaskError(stale)
return guard
def _require_cancel_generation(row: sqlite3.Row, expected_cancel_generation: int) -> None:
if int(row["cancel_generation"]) != expected_cancel_generation:
raise StaleTaskError("task cancellation generation changed")
@@ -560,15 +562,51 @@ def _transition(
def _generation_transition(
db_path: Path | str, identity: TaskIdentity, lease: DriverLease, *, status: str,
execution_generation: int, cancel_generation: int, generation_stale: str, sql: str, set_params: tuple[Any, ...],
stale: str, now: float, replay: Callable[[sqlite3.Row], Any] | None = None,
db_path: Path | str, identity: TaskIdentity, lease: DriverLease, name: str, execution_generation: int,
cancel_generation: int, *, now: float, set_params: tuple[Any, ...],
replay: Callable[[sqlite3.Row], Any] | None = None,
) -> dict[str, Any]:
"""Lease-first transition fenced on ``status`` + both generations (the recovery family)."""
"""Lease-first transition from ``_GENERATION_TRANSITIONS`` fenced on status + both generations."""
status, set_clause, generation_stale, stale = _GENERATION_TRANSITIONS[name]
def guard(row: sqlite3.Row) -> None:
if not _generations_match(row, status, execution_generation, cancel_generation):
raise StaleTaskError(generation_stale)
return _transition(
db_path, identity, lease=lease, now=now, replay=replay,
guard=_generation_guard(status, execution_generation, cancel_generation, generation_stale),
sql=sql, set_params=set_params, fence_params=(execution_generation, cancel_generation), stale=stale,
db_path, identity, lease=lease, now=now, replay=replay, guard=guard,
sql=_generation_update(set_clause, status), set_params=set_params,
fence_params=(execution_generation, cancel_generation), stale=stale,
)
def _run_fence_transition(
db_path: Path | str, attempt: TaskAttempt, *, guard_stale: str,
lease_generation: Callable[[Any], int] = int, **transition: Any,
) -> dict[str, Any]:
"""Transition fenced on this attempt's running generation under its exact lease (row guard + SQL fence).
``lease_generation`` casts the stored run_lease_generation: ``int`` raises on NULL,
``int(v or 0)`` reads it as generation 0.
"""
lease = attempt.lease
def guard(row: sqlite3.Row) -> None:
if not (
_generations_match(row, "running", attempt.execution_generation, attempt.cancel_generation)
and row["run_gateway_id"] == lease.gateway_id
and row["run_process_generation"] == lease.process_generation
and lease_generation(row["run_lease_generation"]) == lease.lease_generation
):
raise StaleTaskError(guard_stale)
return _transition(
db_path, attempt.identity, lease=lease, guard=guard,
fence_params=(
attempt.execution_generation, attempt.cancel_generation,
lease.gateway_id, lease.process_generation, lease.lease_generation,
),
**transition,
)
@@ -592,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())
@@ -756,25 +795,10 @@ def settle_task(
settlement_id = _terminal_settlement_id(settlement_id, status)
result_json = _canonical_json(result)
now = _timestamp(clock)
lease, identity = attempt.lease, attempt.identity
def guard(row: sqlite3.Row) -> None:
if not (
_generations_match(row, "running", attempt.execution_generation, attempt.cancel_generation)
and row["run_gateway_id"] == lease.gateway_id
and row["run_process_generation"] == lease.process_generation
and int(row["run_lease_generation"]) == lease.lease_generation
):
raise StaleTaskError("task attempt is stale or cancelled")
return _transition(
db_path, identity, lease=lease, lease_first=False, now=now,
replay=_settlement_replay(settlement_id, status, result_json), guard=guard,
return _run_fence_transition(
db_path, attempt, guard_stale="task attempt is stale or cancelled", lease_first=False, now=now,
replay=_settlement_replay(settlement_id, status, result_json),
sql=_SETTLE_RUNNING_SQL, set_params=(status, settlement_id, status, result_json, now, now),
fence_params=(
attempt.execution_generation, attempt.cancel_generation,
lease.gateway_id, lease.process_generation, lease.lease_generation,
),
stale="task changed during settlement",
)
@@ -808,11 +832,9 @@ def resolve_indeterminate_task(
result_json = _canonical_json(result)
now = _timestamp(clock)
return _generation_transition(
db_path, identity, lease, status="indeterminate", execution_generation=expected_execution_generation,
cancel_generation=expected_cancel_generation, generation_stale="indeterminate task generation changed",
replay=_settlement_replay(settlement_id, status, result_json), now=now,
sql=_SETTLE_INDETERMINATE_SQL, set_params=(status, settlement_id, status, result_json, now, now),
stale="indeterminate task changed during reconciliation",
db_path, identity, lease, "resolve", expected_execution_generation, expected_cancel_generation, now=now,
replay=_settlement_replay(settlement_id, status, result_json),
set_params=(status, settlement_id, status, result_json, now, now),
)
@@ -825,11 +847,8 @@ def resolve_indeterminate_cancellation(
cancel_id = _identifier(cancel_id, label="cancel_id")
now = _timestamp(clock)
return _generation_transition(
db_path, identity, lease, status="indeterminate", execution_generation=expected_execution_generation,
cancel_generation=expected_cancel_generation, generation_stale="indeterminate cancellation proof is stale",
replay=_cancel_replay(cancel_id), now=now,
sql=_CANCEL_INDETERMINATE_SQL, set_params=(expected_cancel_generation + 1, cancel_id, now, now),
stale="indeterminate cancellation proof lost its fence",
db_path, identity, lease, "resolve_cancel", expected_execution_generation, expected_cancel_generation,
now=now, replay=_cancel_replay(cancel_id), set_params=(expected_cancel_generation + 1, cancel_id, now, now),
)
@@ -841,10 +860,8 @@ def requeue_indeterminate_task(
_expected_generations(lease, identity, expected_execution_generation, expected_cancel_generation)
now = _timestamp(clock)
return _generation_transition(
db_path, identity, lease, status="indeterminate", execution_generation=expected_execution_generation,
cancel_generation=expected_cancel_generation, generation_stale="indeterminate task generation changed",
now=now, sql=_REQUEUE_INDETERMINATE_SQL, set_params=(now,),
stale="indeterminate task changed during requeue",
db_path, identity, lease, "requeue", expected_execution_generation, expected_cancel_generation, now=now,
set_params=(now,),
)
@@ -863,10 +880,8 @@ def defer_indeterminate_task(
return _task_from_row(row, idempotent=True) if deferred and row["result_json"] == result_json else None
return _generation_transition(
db_path, identity, lease, status="indeterminate", execution_generation=expected_execution_generation,
cancel_generation=expected_cancel_generation, generation_stale="indeterminate task generation changed",
replay=replay, now=now, sql=_DEFER_SQL, set_params=(result_json, now, now),
stale="indeterminate task changed during deferral",
db_path, identity, lease, "defer", expected_execution_generation, expected_cancel_generation, now=now,
replay=replay, set_params=(result_json, now, now),
)
@@ -878,18 +893,15 @@ def requeue_deferred_task(
_expected_generations(lease, identity, expected_execution_generation, expected_cancel_generation)
now = _timestamp(clock)
return _generation_transition(
db_path, identity, lease, status="deferred", execution_generation=expected_execution_generation,
cancel_generation=expected_cancel_generation, generation_stale="deferred task generation changed",
now=now, sql=_REQUEUE_DEFERRED_SQL, set_params=(now,),
stale="deferred task changed during requeue",
db_path, identity, lease, "requeue_deferred", expected_execution_generation, expected_cancel_generation,
now=now, set_params=(now,),
)
def requeue_not_admitted_task(db_path: Path | str, attempt: TaskAttempt, *, clock: Clock) -> dict[str, Any]:
"""Return a running task to its durable queue after proven non-admission."""
now = _timestamp(clock)
lease, identity = attempt.lease, attempt.identity
_check_same_room(lease, identity)
_check_same_room(attempt.lease, attempt.identity)
def replay(row: sqlite3.Row) -> dict[str, Any] | None:
requeued = (
@@ -900,23 +912,10 @@ def requeue_not_admitted_task(db_path: Path | str, attempt: TaskAttempt, *, cloc
)
return _task_from_row(row, idempotent=True) if requeued else None
def guard(row: sqlite3.Row) -> None:
if not (
_generations_match(row, "running", attempt.execution_generation, attempt.cancel_generation)
and row["run_gateway_id"] == lease.gateway_id
and row["run_process_generation"] == lease.process_generation
and int(row["run_lease_generation"] or 0) == lease.lease_generation
):
raise StaleTaskError("not-admitted task attempt lost its fence")
return _transition(
db_path, identity, lease=lease, now=now, replay=replay, guard=guard, sql=_REQUEUE_RUNNING_SQL,
set_params=(now,),
fence_params=(
attempt.execution_generation, attempt.cancel_generation,
lease.gateway_id, lease.process_generation, lease.lease_generation,
),
stale="not-admitted task changed during requeue",
return _run_fence_transition(
db_path, attempt, guard_stale="not-admitted task attempt lost its fence",
lease_generation=lambda value: int(value or 0), now=now, replay=replay,
sql=_REQUEUE_RUNNING_SQL, set_params=(now,), stale="not-admitted task changed during requeue",
)
@@ -1004,25 +1003,25 @@ def recover_room(db_path: Path | str, lease: DriverLease, *, clock: Clock) -> di
}
def get_task(db_path: Path | str, identity: TaskIdentity) -> dict[str, Any]:
"""Read one task without mutating its state."""
def _read(db_path: Path | str, query: Callable[[sqlite3.Connection], Any]) -> Any:
conn = _connect(db_path)
try:
return _task_from_row(_load_task(conn, identity))
return query(conn)
finally:
conn.close()
def get_task(db_path: Path | str, identity: TaskIdentity) -> dict[str, Any]:
"""Read one task without mutating its state."""
return _read(db_path, lambda conn: _task_from_row(_load_task(conn, identity)))
def list_tasks(db_path: Path | str, *, room_id: Any, status: TaskStatus | None = None) -> list[dict[str, Any]]:
"""Return room tasks in deterministic admission order."""
room_id = _identifier(room_id, label="room_id")
if status is not None and status not in TASK_STATUSES:
raise DriverValidationError("invalid task status")
conn = _connect(db_path)
try:
return [_task_from_row(row) for row in _tasks_in_order(conn, room_id, status)]
finally:
conn.close()
return _read(db_path, lambda conn: [_task_from_row(row) for row in _tasks_in_order(conn, room_id, status)])
def prune_published_terminal_tasks(
@@ -1043,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()

View File

@@ -19,7 +19,7 @@ import stat
import time
import urllib.parse
from dataclasses import asdict, dataclass
from functools import lru_cache
from functools import lru_cache, partial
from pathlib import Path
from typing import Any, Callable, Iterable, Literal, Mapping
@@ -146,13 +146,10 @@ def _digest(value: Any, *, field: str) -> str:
return value
def _exact_fields(
value: Mapping[str, Any], *, required: set[str], optional: set[str] = frozenset(), label: str
) -> None:
exact_fields(
value, label=label, required=required, optional=optional, error=HostedRoomPeerError,
missing_fmt="{label} missing fields: {fields}", unknown_fmt="{label} unknown fields: {fields}",
)
_exact_fields = partial(
exact_fields, error=HostedRoomPeerError, missing_fmt="{label} missing fields: {fields}",
unknown_fmt="{label} unknown fields: {fields}",
)
def _canonical_json(value: Mapping[str, Any]) -> bytes:

View File

@@ -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 (
@@ -92,6 +93,20 @@ def _require_room(conn: sqlite3.Connection, room_id: str) -> None:
raise hosted_rooms.RoomNotFoundError("hosted room not found")
def _settled_message(
conn: sqlite3.Connection, room_id: str, discussion_event_id: str, message_event_id: Any
) -> dict[str, Any] | None:
"""Return the indexed member message a ``turn.settled`` event committed, if it is in the projection."""
rows = conn.execute(
"SELECT seq, event_json FROM hosted_room_policy_events WHERE room_id=? AND discussion_event_id=?",
(room_id, discussion_event_id),
).fetchall()
return next(
(m for m in (json.loads(row["event_json"]) for row in rows) if m.get("event_id") == message_event_id),
None,
)
class HostedRoomPolicyCheckpoint:
"""Incrementally index room policy without compacting visible history."""
@@ -115,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=(",", ":")),
@@ -129,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"}:
@@ -213,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)
@@ -249,29 +267,18 @@ 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 "")
seen_through_seq = int(payload.get("seen_through_seq") or 0)
if kind == "turn.settled" and payload.get("message_event_id"):
messages = conn.execute(
"SELECT seq, event_json FROM hosted_room_policy_events WHERE room_id=? AND discussion_event_id=?",
(room_id, discussion_event_id),
).fetchall()
committed = next(
(
message for message in (json.loads(row["event_json"]) for row in messages)
if message.get("event_id") == payload["message_event_id"]
),
None,
)
committed = _settled_message(conn, room_id, discussion_event_id, payload["message_event_id"])
if committed is not None:
seen_through_seq = max(seen_through_seq, int(committed["seq"]))
self._store_transcript_event(
conn, event=committed, thread_id=thread_id, settled_seq=int(event["seq"])
)
self._store_transcript_event(conn, event=committed, thread_id=thread_id, settled_seq=seq)
if member_id and seen_through_seq > 0:
conn.execute(
"""INSERT INTO hosted_room_policy_watermarks(
@@ -322,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(

View File

@@ -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)
)"""
)
@@ -148,6 +152,43 @@ def _validate_page(page: Any) -> tuple[list[dict[str, Any]], dict[str, Any]]:
return events, {"gateway_id": gateway_id, "epoch": epoch}
def _replica_row_state(conn: sqlite3.Connection, room_id: str) -> tuple[sqlite3.Row | None, int, int, int]:
"""Return (row, stored_epoch, last_seq, stored_bytes); a new room is admitted only under the room cap."""
row = conn.execute(
"""SELECT authority_gateway_id, authority_epoch, last_seq, latest_seq, event_bytes
FROM hosted_room_replicas WHERE room_id=?""",
(room_id,),
).fetchone()
if row is None:
count = conn.execute("SELECT COUNT(*) FROM hosted_room_replicas").fetchone()[0]
if int(count) >= MAX_REPLICA_ROOMS:
raise ReplicaError("replica room capacity exhausted")
return None, 0, 0, 0
return row, int(row["authority_epoch"]), int(row["last_seq"]), int(row["event_bytes"])
def _store_replica(
conn: sqlite3.Connection, *, is_new: bool, room_id: str, room_name: str, members_json: str,
authority: dict[str, Any], new_last: int, latest_seq: int, added_bytes: int, now: float,
) -> None:
"""INSERT the replica row for a new room, else UPDATE it (event_bytes accumulates)."""
values = (room_name, members_json, authority["gateway_id"], authority["epoch"], new_last, max(latest_seq, new_last))
if is_new:
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 (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(room_id, *values, added_bytes, now, now),
)
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=?""",
(*values, added_bytes, now, room_id),
)
def ingest_page(
db_path: Path | str, *, room_id: Any, room_name: Any, members: Any, page: Any, now: float | None = None
) -> dict[str, Any]:
@@ -163,20 +204,7 @@ def ingest_page(
now = time.time() if now is None else float(now)
with _replica_transaction(db_path) as conn:
row = conn.execute(
"""SELECT authority_gateway_id, authority_epoch, last_seq, latest_seq, event_bytes
FROM hosted_room_replicas WHERE room_id=?""",
(room_id,),
).fetchone()
if row is None:
count = conn.execute("SELECT COUNT(*) FROM hosted_room_replicas").fetchone()[0]
if int(count) >= MAX_REPLICA_ROOMS:
raise ReplicaError("replica room capacity exhausted")
stored_epoch = last_seq = stored_bytes = 0
else:
stored_epoch, last_seq, stored_bytes = (
int(row["authority_epoch"]), int(row["last_seq"]), int(row["event_bytes"])
)
row, stored_epoch, last_seq, stored_bytes = _replica_row_state(conn, room_id)
if authority["epoch"] < stored_epoch:
raise ReplicaEpochRegressionError("page authority epoch is older than the stored replica epoch")
@@ -202,26 +230,10 @@ def ingest_page(
latest_seq = page.get("latest_seq")
if isinstance(latest_seq, bool) or not isinstance(latest_seq, int):
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 (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(
room_id, room_name, members_json, authority["gateway_id"], authority["epoch"],
new_last, max(latest_seq, new_last), added_bytes, now, now,
),
)
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=?""",
(
room_name, members_json, authority["gateway_id"], authority["epoch"],
new_last, max(latest_seq, new_last), added_bytes, now, room_id,
),
)
_store_replica(
conn, is_new=row is None, room_id=room_id, room_name=room_name, members_json=members_json,
authority=authority, new_last=new_last, latest_seq=latest_seq, added_bytes=added_bytes, now=now,
)
return {
"room_id": room_id,
"stored_seq": new_last,
@@ -293,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,
@@ -379,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 {

View File

@@ -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 = (
@@ -265,6 +298,12 @@ def _bounded_limit(value: Any, maximum: int) -> int:
return value
def _non_negative(value: Any, label: str) -> int:
return non_negative_int(
value, error=HostedRoomError, message=f"{label} must be a non-negative integer"
)
def _now(now: float | None) -> float:
return time.time() if now is None else float(now)
@@ -384,13 +423,17 @@ def _primary_key_columns(conn: sqlite3.Connection, table: str) -> tuple[str, ...
return tuple(str(row[1]) for row in sorted(rows, key=lambda row: int(row[5])))
def _remote_run_schema_current(conn: sqlite3.Connection, columns: frozenset[str]) -> bool:
return (
_REMOTE_RUN_SCHEMA_COLUMNS.issubset(columns)
and _primary_key_columns(conn, "hosted_room_remote_runs") == _REMOTE_RUN_IDENTITY_COLUMNS
)
def _migrate_remote_run_schema(conn: sqlite3.Connection) -> None:
"""Fence legacy receipts behind a complete authority-lineage key."""
columns = table_columns(conn, "hosted_room_remote_runs")
if (
_REMOTE_RUN_SCHEMA_COLUMNS.issubset(columns)
and _primary_key_columns(conn, "hosted_room_remote_runs") == _REMOTE_RUN_IDENTITY_COLUMNS
):
if _remote_run_schema_current(conn, columns):
return
conn.execute("DROP TABLE IF EXISTS hosted_room_remote_runs_migrating")
conn.execute(f"CREATE TABLE hosted_room_remote_runs_migrating ({_REMOTE_RUNS_BODY})")
@@ -414,39 +457,57 @@ def _migrate_remote_run_schema(conn: sqlite3.Connection) -> None:
conn.execute("ALTER TABLE hosted_room_remote_runs_migrating RENAME TO hosted_room_remote_runs")
# Draft builds before the actor contract carried no identity. Preserve their
# inert replay rows explicitly as legacy system events rather than guessing a
# user or Bot author.
_LEGACY_ACTOR_JSON = _system_actor_json("legacy").replace("'", "''")
# (table, column, ddl) applied in this exact order; each table's PRAGMA is read on first use.
_LEGACY_COLUMN_DDL = (
(
"hosted_rooms", "authority_gateway_id",
"ALTER TABLE hosted_rooms ADD COLUMN authority_gateway_id TEXT NOT NULL DEFAULT 'legacy'",
),
(
"hosted_rooms", "authority_epoch",
"ALTER TABLE hosted_rooms ADD COLUMN authority_epoch INTEGER NOT NULL DEFAULT 1",
),
(
"hosted_rooms", "event_bytes",
"ALTER TABLE hosted_rooms ADD COLUMN event_bytes INTEGER NOT NULL DEFAULT 0",
),
(
"hosted_room_events", "actor_json",
"ALTER TABLE hosted_room_events "
f"ADD COLUMN actor_json TEXT NOT NULL DEFAULT '{_LEGACY_ACTOR_JSON}'",
),
(
"hosted_room_events", "authority_epoch",
"ALTER TABLE hosted_room_events ADD COLUMN authority_epoch INTEGER",
),
)
def _migrate_legacy_columns(conn: sqlite3.Connection) -> None:
"""Add columns draft schemas lacked; backfill event_bytes when first introduced."""
room_columns = table_columns(conn, "hosted_rooms")
if "authority_gateway_id" not in room_columns:
columns: dict[str, frozenset[str]] = {}
for table, column, ddl in _LEGACY_COLUMN_DDL:
if table not in columns:
columns[table] = table_columns(conn, table)
if column not in columns[table]:
conn.execute(ddl)
if "event_bytes" not in columns["hosted_rooms"]:
conn.execute(
"ALTER TABLE hosted_rooms "
"ADD COLUMN authority_gateway_id TEXT NOT NULL DEFAULT 'legacy'"
)
if "authority_epoch" not in room_columns:
conn.execute(
"ALTER TABLE hosted_rooms ADD COLUMN authority_epoch INTEGER NOT NULL DEFAULT 1"
)
backfill_event_bytes = "event_bytes" not in room_columns
if backfill_event_bytes:
conn.execute("ALTER TABLE hosted_rooms ADD COLUMN event_bytes INTEGER NOT NULL DEFAULT 0")
event_columns = table_columns(conn, "hosted_room_events")
if "actor_json" not in event_columns:
# Draft builds before the actor contract carried no identity. Preserve
# their inert replay rows explicitly as legacy system events rather
# than guessing a user or Bot author.
escaped_actor = _system_actor_json("legacy").replace("'", "''")
conn.execute(
"ALTER TABLE hosted_room_events "
f"ADD COLUMN actor_json TEXT NOT NULL DEFAULT '{escaped_actor}'"
)
if "authority_epoch" not in event_columns:
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)"""
)
@@ -474,10 +535,7 @@ def _schema_is_current(conn: sqlite3.Connection) -> bool:
for (table, required), columns in zip(_REQUIRED_COLUMNS, actual, strict=True):
if not required.issubset(columns):
return False
if (
table == "hosted_room_remote_runs"
and _primary_key_columns(conn, table) != _REMOTE_RUN_IDENTITY_COLUMNS
):
if table == "hosted_room_remote_runs" and not _remote_run_schema_current(conn, columns):
return False
index = conn.execute(
"SELECT 1 FROM sqlite_master WHERE type='index' AND name='idx_hosted_room_events_cursor'"
@@ -576,6 +634,14 @@ def _raise_room_not_found(conn: sqlite3.Connection, room_id: str) -> NoReturn:
raise RoomNotFoundError("hosted room not found")
def _reload(conn: sqlite3.Connection, sql: str, params: tuple, missing: str) -> sqlite3.Row:
"""Re-read a row this transaction just wrote; a miss is an invariant violation."""
row = conn.execute(sql, params).fetchone()
if row is None: # pragma: no cover - guarded by the write above
raise RuntimeError(missing)
return row
def _room_from_row(row: sqlite3.Row, *, idempotent: bool = False) -> dict[str, Any]:
room = {
"room_id": row["room_id"],
@@ -615,18 +681,16 @@ def _load_event(conn: sqlite3.Connection, room_id: str, event_id: str) -> sqlite
return conn.execute(_SELECT_EVENT, (room_id, event_id)).fetchone()
def _event_storage_bytes(event_id: str, kind: str, actor_json: str, payload_json: str) -> int:
return len((event_id + kind + actor_json + payload_json).encode("utf-8"))
def _gateway_event_bytes(conn: sqlite3.Connection) -> int:
return int(conn.execute(_SUM_EVENT_BYTES).fetchone()[0])
def _assert_event_capacity(
conn: sqlite3.Connection, *, room: sqlite3.Row, additional_bytes: int,
allow_control: bool = False,
) -> None:
def _prepare_event(
conn: sqlite3.Connection, room: sqlite3.Row, event_id: str, kind: str, actor_json: str,
payload_json: str, *, allow_control: bool = False,
) -> int:
"""Size one pending event and enforce per-room and gateway capacity; returns its bytes."""
additional_bytes = len((event_id + kind + actor_json + payload_json).encode("utf-8"))
count_reserve = CONTROL_EVENT_COUNT_RESERVE if allow_control else 0
byte_reserve = CONTROL_EVENT_BYTE_RESERVE if allow_control else 0
gateway_byte_limit = MAX_GATEWAY_EVENT_BYTES + byte_reserve
@@ -648,18 +712,7 @@ def _assert_event_capacity(
raise HostedRoomError(
"Group Chat storage is full on this host. Delete an old Group Chat and try again."
)
def _prepare_event(
conn: sqlite3.Connection, room: sqlite3.Row, event_id: str, kind: str, actor_json: str,
payload_json: str, *, allow_control: bool = False,
) -> int:
"""Size one pending event and enforce capacity; returns its storage bytes."""
event_bytes = _event_storage_bytes(event_id, kind, actor_json, payload_json)
_assert_event_capacity(
conn, room=room, additional_bytes=event_bytes, allow_control=allow_control
)
return event_bytes
return additional_bytes
# --- retention -------------------------------------------------------------------
@@ -735,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]
@@ -756,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"],
@@ -820,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(
@@ -895,29 +959,33 @@ 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),
)
def _read_one(db_path: Path | str, sql: str, params: tuple[Any, ...]) -> sqlite3.Row | None:
with _transaction(db_path) as conn:
return conn.execute(sql, params).fetchone()
def peer_room_is_reserved(
db_path: Path | str, *, room_id: str, target_profile: str, now: float | None = None
) -> bool:
"""Return whether a live target-side RoomLink reservation fences Desktop."""
timestamp = _now(now)
with _transaction(db_path) as conn:
row = conn.execute(
_SELECT_LIVE_RESERVATION,
(_room_id(room_id), _actor_id(target_profile, "target_profile"), timestamp),
).fetchone()
return row is not None
params = (_room_id(room_id), _actor_id(target_profile, "target_profile"), timestamp)
return _read_one(db_path, _SELECT_LIVE_RESERVATION, params) is not None
def peer_room_grant_is_current(
@@ -926,13 +994,13 @@ def peer_room_grant_is_current(
"""Require a grant to match the target's current live reservation."""
timestamp = _now(now)
values = _reservation_claims(claims)
with _transaction(db_path) as conn:
row = conn.execute(
"""SELECT 1 FROM hosted_room_peer_reservations WHERE room_id=? AND member_id=?
AND target_profile=? AND authority_gateway_id=? AND authority_epoch=?
AND expires_at>? AND revoked_at IS NULL LIMIT 1""",
(*values, timestamp),
).fetchone()
row = _read_one(
db_path,
"""SELECT 1 FROM hosted_room_peer_reservations WHERE room_id=? AND member_id=?
AND target_profile=? AND authority_gateway_id=? AND authority_epoch=?
AND expires_at>? AND revoked_at IS NULL LIMIT 1""",
(*values, timestamp),
)
return row is not None
@@ -943,12 +1011,12 @@ def room_grant_is_revoked(
timestamp = _now(now)
scope_key = _room_grant_scope_key(claims)
issued_at = float(claims.get("issued_at") or 0)
with _transaction(db_path) as conn:
row = conn.execute(
"""SELECT revoked_before FROM hosted_room_revoked_grants
WHERE scope_key=? AND expires_at>?""",
(scope_key, timestamp),
).fetchone()
row = _read_one(
db_path,
"""SELECT revoked_before FROM hosted_room_revoked_grants
WHERE scope_key=? AND expires_at>?""",
(scope_key, timestamp),
)
return row is not None and issued_at <= float(row["revoked_before"])
@@ -977,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),
)
@@ -1009,11 +1079,11 @@ def list_remote_run_receipts(
def remote_run_receipt(db_path: Path | str, *, record: Mapping[str, Any]) -> dict[str, Any] | None:
"""Return the exact durable remote run handle for one task attempt."""
with _transaction(db_path) as conn:
row = conn.execute(
f"SELECT * FROM hosted_room_remote_runs WHERE {_REMOTE_RUN_WHERE}",
_remote_run_identity(record),
).fetchone()
row = _read_one(
db_path,
f"SELECT * FROM hosted_room_remote_runs WHERE {_REMOTE_RUN_WHERE}",
_remote_run_identity(record),
)
return dict(row) if row is not None else None
@@ -1041,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""",
(
@@ -1052,14 +1123,13 @@ def _adopt_legacy_room(
)
if adopted.rowcount != 1:
raise AuthorityConflictError("legacy room adoption lost its fence")
existing = conn.execute(_SELECT_ROOM, (room_id,)).fetchone()
if existing is None: # pragma: no cover - row updated above
raise RuntimeError("adopted room could not be reloaded")
existing = _reload(conn, _SELECT_ROOM, (room_id,), "adopted room could not be reloaded")
result = _room_from_row(existing, idempotent=True)
result["adopted"] = True
claim_event = _load_event(conn, room_id, "system:authority-adopted")
if claim_event is None: # pragma: no cover - inserted above
raise RuntimeError("legacy adoption event could not be reloaded")
claim_event = _reload(
conn, _SELECT_EVENT, (room_id, "system:authority-adopted"),
"legacy adoption event could not be reloaded",
)
result["claim_event"] = _event_from_row(claim_event)
return result
@@ -1116,13 +1186,13 @@ def create_room(
VALUES (?, ?, ?, ?, 1, 1, 0, 1, ?, ?, NULL)""",
(room_id, name, members_json, authority_gateway_id, now, now),
)
row = conn.execute(
row = _reload(
conn,
"""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
raise RuntimeError("created room could not be reloaded")
"created room could not be reloaded",
)
result = _room_from_row(row)
result["members"] = normalized_members
return result
@@ -1134,9 +1204,7 @@ def list_rooms(
) -> list[dict[str, Any]]:
"""Return one bounded read-only page ordered by most recent change."""
limit = _bounded_limit(limit, MAX_ROOM_LIST_LIMIT)
offset = non_negative_int(
offset, error=HostedRoomError, message="offset must be a non-negative integer"
)
offset = _non_negative(offset, "offset")
conn = _read_connection(db_path)
try:
rows = conn.execute(
@@ -1177,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(
@@ -1254,12 +1323,12 @@ def append_event(
)
if advanced.rowcount != 1:
raise RuntimeError("hosted room sequence advance lost its write fence")
row = conn.execute(
row = _reload(
conn,
f"SELECT {_EVENT_COLUMNS} FROM hosted_room_events WHERE room_id=? AND seq=?",
(room_id, seq),
).fetchone()
if row is None: # pragma: no cover - guarded by the insert above
raise RuntimeError("appended event could not be reloaded")
"appended event could not be reloaded",
)
result = _event_from_row(row)
result["actor"] = normalized_actor
return result
@@ -1354,6 +1423,32 @@ def request_room_stop(
)
def _append_authority_claim(
conn: sqlite3.Connection, row: sqlite3.Row, *, room_id: str, event_id: str,
expected_gateway_id: str, expected_epoch: int, new_gateway_id: str, target_epoch: int,
actor_json: str, payload_json: str, now: float,
) -> sqlite3.Row | None:
"""Insert the claim event and CAS the room's authority; returns the stored claim event."""
seq = int(row["next_seq"])
claim_bytes = _prepare_event(
conn, row, event_id, "authority.claimed", actor_json, payload_json, allow_control=True
)
conn.execute(
_INSERT_EVENT,
(room_id, seq, event_id, "authority.claimed", actor_json, target_epoch, payload_json, now),
)
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=?""",
(new_gateway_id, claim_bytes, now, room_id, expected_gateway_id, expected_epoch),
)
if updated.rowcount != 1:
raise AuthorityConflictError("hosted room authority changed")
return _load_event(conn, room_id, event_id)
def claim_authority(
db_path: Path | str, *, room_id: Any, expected_gateway_id: Any, expected_epoch: Any,
new_gateway_id: Any, event_id: Any, now: float | None = None,
@@ -1385,7 +1480,8 @@ def claim_authority(
current_gateway = str(row["authority_gateway_id"])
current_epoch = int(row["authority_epoch"])
existing_event = _load_event(conn, room_id, event_id)
if existing_event is not None:
idempotent = existing_event is not None
if idempotent:
if (
existing_event["kind"] != "authority.claimed"
or existing_event["actor_json"] != claim_actor_json
@@ -1395,40 +1491,22 @@ def claim_authority(
raise EventConflictError("event_id already exists with different content")
if current_gateway != new_gateway_id or current_epoch != target_epoch:
raise AuthoritySupersededError("authority claim succeeded but was later superseded")
idempotent = True
elif current_gateway != expected_gateway_id or current_epoch != expected_epoch:
raise AuthorityConflictError("hosted room authority changed")
else:
seq = int(row["next_seq"])
claim_bytes = _prepare_event(
conn, row, event_id, "authority.claimed", claim_actor_json, claim_payload_json,
allow_control=True,
existing_event = _append_authority_claim(
conn, row, room_id=room_id, event_id=event_id,
expected_gateway_id=expected_gateway_id, expected_epoch=expected_epoch,
new_gateway_id=new_gateway_id, target_epoch=target_epoch,
actor_json=claim_actor_json, payload_json=claim_payload_json, now=now,
)
conn.execute(
_INSERT_EVENT,
(
room_id, seq, event_id, "authority.claimed", claim_actor_json, target_epoch,
claim_payload_json, now,
),
)
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=?""",
(new_gateway_id, claim_bytes, now, room_id, expected_gateway_id, expected_epoch),
)
if updated.rowcount != 1:
raise AuthorityConflictError("hosted room authority changed")
idempotent = False
existing_event = _load_event(conn, room_id, event_id)
state_row = conn.execute(
state_row = _reload(
conn,
"""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
raise RuntimeError("claimed room could not be reloaded")
"claimed room could not be reloaded",
)
state = _room_from_row(state_row, idempotent=idempotent)
state["latest_seq"] = int(state_row["next_seq"]) - 1
if existing_event is None: # pragma: no cover - both claim paths set it
@@ -1437,6 +1515,30 @@ def claim_authority(
return state
def _disband_replay(
conn: sqlite3.Connection, room_id: str, room: sqlite3.Row | None
) -> dict[str, Any] | None:
"""Idempotent replay for a retired or already-disbanded room; None when the room is live."""
if room is None:
retired = conn.execute(
"SELECT retired_at FROM hosted_room_retired_ids WHERE room_id=?", (room_id,)
).fetchone()
if retired is None:
raise RoomNotFoundError("hosted room not found")
return {
"room_id": room_id, "disbanded_at": float(retired["retired_at"]),
"idempotent": True, "history_expired": True,
}
if room["disbanded_at"] is None:
return None
conn.execute(_INSERT_RETIRED, (room_id, float(room["disbanded_at"])))
event = _load_event(conn, room_id, "system:room-disbanded")
result = {"room_id": room_id, "disbanded_at": float(room["disbanded_at"]), "idempotent": True}
if event is not None:
result["event"] = _event_from_row(event, idempotent=True)
return result
def disband_room(
db_path: Path | str, *, room_id: Any, expected_gateway_id: Any, expected_epoch: Any,
now: float | None = None,
@@ -1453,25 +1555,9 @@ def disband_room(
FROM hosted_rooms WHERE room_id=?""",
(room_id,),
).fetchone()
if room is None:
retired = conn.execute(
"SELECT retired_at FROM hosted_room_retired_ids WHERE room_id=?", (room_id,)
).fetchone()
if retired is None:
raise RoomNotFoundError("hosted room not found")
return {
"room_id": room_id, "disbanded_at": float(retired["retired_at"]),
"idempotent": True, "history_expired": True,
}
if room["disbanded_at"] is not None:
conn.execute(_INSERT_RETIRED, (room_id, float(room["disbanded_at"])))
event = _load_event(conn, room_id, "system:room-disbanded")
result = {
"room_id": room_id, "disbanded_at": float(room["disbanded_at"]), "idempotent": True
}
if event is not None:
result["event"] = _event_from_row(event, idempotent=True)
return result
replay = _disband_replay(conn, room_id, room)
if replay is not None:
return replay
if (
str(room["authority_gateway_id"]) != expected_gateway_id
or int(room["authority_epoch"]) != expected_epoch
@@ -1492,17 +1578,20 @@ 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:
raise RoomConflictError("hosted room disband lost its fence")
conn.execute(_INSERT_RETIRED, (room_id, now))
event = _load_event(conn, room_id, "system:room-disbanded")
if event is None: # pragma: no cover - inserted in this transaction
raise RuntimeError("room disband event could not be reloaded")
event = _reload(
conn, _SELECT_EVENT, (room_id, "system:room-disbanded"),
"room disband event could not be reloaded",
)
_prune_disbanded_rooms_locked(
conn, now=now, max_gateway_event_bytes=MAX_GATEWAY_EVENT_BYTES
)
@@ -1518,9 +1607,7 @@ def read_events(
) -> dict[str, Any]:
"""Read a monotonic room-log delta after ``since_seq``."""
room_id = _room_id(room_id)
since_seq = non_negative_int(
since_seq, error=HostedRoomError, message="since_seq must be a non-negative integer"
)
since_seq = _non_negative(since_seq, "since_seq")
limit = _bounded_limit(limit, MAX_LOG_LIMIT)
with _transaction(db_path) as conn: