diff --git a/gateway/hosted_room_discussion.py b/gateway/hosted_room_discussion.py index 1be91f6866..0dc4e9b225 100644 --- a/gateway/hosted_room_discussion.py +++ b/gateway/hosted_room_discussion.py @@ -1,10 +1,9 @@ """Deterministic policy for same-gateway hosted-room Discussions. -Pure (no I/O, transport or model knowledge): a frozen local roster plus the complete -typed room log yields at most one next driver task. Discussion coordinates live in -deterministic ``TaskIdentity`` values and typed terminal events rather than a widened -driver payload, so a restart reconstructs tasks from durable state. Callers must -reconcile terminal driver rows into publication plans before asking for the next task. +Pure (no I/O, transport or model knowledge): a frozen local roster plus the complete typed room log yields +at most one next driver task. Discussion coordinates live in deterministic ``TaskIdentity`` values and typed +terminal events rather than a widened driver payload, so a restart reconstructs tasks from durable state. +Callers must reconcile terminal driver rows into publication plans before asking for the next task. """ from __future__ import annotations @@ -754,9 +753,8 @@ def plan_publication( result: Any = None, execution_generation: int | None = None, local_profiles: Iterable[str]) -> PublicationPlan: """Plan idempotent room effects for one terminal driver task. - A newer user event in the same thread supersedes a late result: the task - stays terminal in driver state, but only a deterministic cancellation is - published so stale prose and its watermark cannot hide the newer message. + A newer user event in the same thread supersedes a late result: the task stays terminal in driver state, + but only a deterministic cancellation is published so stale prose and its watermark cannot hide it. """ room = validate_room(room_value, local_profiles=local_profiles) validated = _validated_events(events, room=room) diff --git a/gateway/hosted_room_driver.py b/gateway/hosted_room_driver.py index b53c85909a..fb7268767a 100644 --- a/gateway/hosted_room_driver.py +++ b/gateway/hosted_room_driver.py @@ -1,9 +1,8 @@ """Durable execution state for a same-gateway hosted room driver. -Owns only the driver lease and task state machine: no model calls, no -sessions, no dependency on the hosted-room event log. Callers supply the -database path and clock so recovery and fencing are testable without -process-global state. +Owns only the driver lease and task state machine: no model calls, no sessions, no dependency on the +hosted-room event log. Callers supply the database path and clock so recovery and fencing are testable +without process-global state. """ from __future__ import annotations @@ -53,8 +52,8 @@ _TASK_INDEX_SQL = """CREATE INDEX {if_not_exists}idx_hosted_room_driver_tasks_st ON hosted_room_driver_tasks(room_id, status, source_event_seq, created_at, task_id)""" # --- Fenced task UPDATE statements (one per state-machine transition) --------- -# Every transition is "UPDATE ... SET WHERE room_id=? AND task_id=? AND "; -# the fence names the expected status plus the generations that must not have moved. +# Every transition is "UPDATE ... SET WHERE room_id=? AND task_id=? AND "; the fence names the +# expected status plus the generations that must not have moved. _GENERATION_FENCE = "execution_generation=? AND cancel_generation=?" _RUN_FENCE = "run_gateway_id=? AND run_process_generation=? AND run_lease_generation=?" _SETTLE_SET = "status=?, settlement_id=?, settlement_status=?, result_json=?, terminal_at=?, updated_at=?" @@ -263,9 +262,7 @@ def _migrate_task_status_constraint(conn: sqlite3.Connection) -> None: def _connect(db_path: DbPath) -> sqlite3.Connection: """Open the store; existing tables are validated (after the status-constraint migration when needed). - - The driver schema never shipped, so an incompatible draft fails closed in ``_validate_schema``. - """ + The driver schema never shipped, so an incompatible draft fails closed in ``_validate_schema``.""" existing: list[bool] = [] def ready(conn: sqlite3.Connection) -> bool: existing.append(_schema_objects_exist(conn)) @@ -417,10 +414,8 @@ def _terminal_settlement_id(settlement_id: Any, status: Any) -> str: def _settlement( settlement_id: Any, status: Any, result: Any, clock: Clock ) -> tuple[float, Callable[[sqlite3.Row], Any], tuple[Any, ...]]: - """Validate one terminal settlement -> (now, replay predicate, ``_SETTLE_SET`` params). - - Replay: an identical settlement already committed is idempotent; a different one is a conflict. - """ + """Validate one terminal settlement -> (now, replay predicate, ``_SETTLE_SET`` params); the replay treats an + identical committed settlement as idempotent and a different one as a conflict.""" settlement_id = _terminal_settlement_id(settlement_id, status) result_json = _canonical_json(result) now = _timestamp(clock) @@ -456,9 +451,9 @@ def _transition( guard: Callable[[sqlite3.Row], None] | None = None) -> dict[str, Any]: """Run one fenced task transition: load -> idempotent replay -> lease/fence guard -> UPDATE. - ``sql`` binds ``(*set_params, room_id, task_id, *fence_params)`` and must hit exactly one row or - ``stale`` is raised. ``lease_first`` checks the lease before the row load (recovery paths) instead - of after the replay (settlement paths: an identical replay still succeeds after the lease moved on). + ``sql`` binds ``(*set_params, room_id, task_id, *fence_params)`` and must hit exactly one row or ``stale`` + is raised. ``lease_first`` checks the lease before the row load (recovery paths) instead of after the + replay (settlement paths: an identical replay still succeeds after the lease moved on). """ params = (*set_params, identity.room_id, identity.task_id, *fence_params) with _transaction(db_path) as conn: @@ -493,9 +488,7 @@ def _run_fence_transition( db_path: DbPath, 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 0. - """ + ``lease_generation`` casts the stored run_lease_generation: ``int`` raises on NULL, ``int(v or 0)`` reads 0.""" lease = attempt.lease def guard(row: sqlite3.Row) -> None: if not _generations_match(row, "running", attempt.execution_generation, attempt.cancel_generation) or ( diff --git a/gateway/hosted_room_links.py b/gateway/hosted_room_links.py index ac24473ce0..000986b20c 100644 --- a/gateway/hosted_room_links.py +++ b/gateway/hosted_room_links.py @@ -1,8 +1,7 @@ """Private SQLite storage for negotiated hosted-room links. -Route metadata and its scoped grant share the gateway's private root -``state.db``. SQLite WAL plus ``BEGIN IMMEDIATE`` owns concurrency; grants are -never included in reprs, status payloads, or exception messages. +Route metadata and its scoped grant share the gateway's private root ``state.db``. SQLite WAL plus +``BEGIN IMMEDIATE`` owns concurrency; grants are never included in reprs, status payloads, or exception messages. """ from __future__ import annotations diff --git a/gateway/hosted_room_peer.py b/gateway/hosted_room_peer.py index fa3cff0a38..4264ea0ec1 100644 --- a/gateway/hosted_room_peer.py +++ b/gateway/hosted_room_peer.py @@ -1,8 +1,7 @@ """Typed contracts for autonomous cross-gateway hosted-room members. -The Desktop may bootstrap an invitation but is never the issuer or runtime -courier: the target gateway verifies a scoped grant and the full task -coordinates before admitting any model or tool work. +The Desktop may bootstrap an invitation but is never the issuer or runtime courier: the target gateway +verifies a scoped grant and the full task coordinates before admitting any model or tool work. """ from __future__ import annotations diff --git a/gateway/hosted_room_policy_checkpoint.py b/gateway/hosted_room_policy_checkpoint.py index c72daafa95..9def25f146 100644 --- a/gateway/hosted_room_policy_checkpoint.py +++ b/gateway/hosted_room_policy_checkpoint.py @@ -1,8 +1,8 @@ """Durable bounded policy projection for hosted Group Chat preparation. -The append-only room log remains the user-visible source of truth. This module -materializes only the state needed to choose and reconstruct the next active -discussion, so a busy room does not replay its complete history every poll. +The append-only room log remains the user-visible source of truth. This module materializes only the state +needed to choose and reconstruct the next active discussion, so a busy room does not replay its complete +history every poll. """ from __future__ import annotations diff --git a/gateway/hosted_room_replicas.py b/gateway/hosted_room_replicas.py index b180de21c9..3a786b2f31 100644 --- a/gateway/hosted_room_replicas.py +++ b/gateway/hosted_room_replicas.py @@ -53,10 +53,8 @@ def _initialize_replica_schema(conn: sqlite3.Connection) -> None: @contextmanager def _replica_transaction(db_path: DbPath) -> Iterator[sqlite3.Connection]: - """Ensure the replica schema (own autocommit connection), then open an IMMEDIATE transaction. - - The DDL is deliberately re-run inside the transaction: that double init is the established statement order. - """ + """Ensure the replica schema (own autocommit connection), then open an IMMEDIATE transaction. The DDL is + deliberately re-run inside the transaction: that double init is the established statement order.""" with closing(_connect(db_path)) as conn, conn: _initialize_replica_schema(conn) with _transaction(db_path, immediate=True) as conn: @@ -202,9 +200,9 @@ def promote_replica( ) -> dict[str, Any]: """Continue a replicated room on THIS gateway at ``epoch + 1``. - Copies the replica log into the authoritative store and appends a lineage-proving - ``authority.claimed`` event, so wherever the claim replicates the old epoch is stale and every - fenced primitive rejects it. The caller decides takeover is safe; this makes it atomic and provable. + Copies the replica log into the authoritative store and appends a lineage-proving ``authority.claimed`` + event, so wherever the claim replicates the old epoch is stale and every fenced primitive rejects it. + The caller decides takeover is safe; this makes it atomic and provable. """ room_id = _room_id(room_id) if not isinstance(reason, str) or not reason or len(reason) > 200: @@ -251,9 +249,9 @@ def demote_room( ) -> dict[str, Any]: """Fence THIS gateway's stale room authority against a proven newer epoch. - When a returning gateway observes (replicated ``authority.claimed`` or a transport rejection) that - another gateway owns the room at a higher epoch, append ``authority.lost`` and adopt the observed - lineage so no local send can commit at the stale epoch. Idempotent per lineage. + When a returning gateway observes (replicated ``authority.claimed`` or a transport rejection) that another + gateway owns the room at a higher epoch, append ``authority.lost`` and adopt the observed lineage so no + local send can commit at the stale epoch. Idempotent per lineage. """ room_id = _room_id(room_id) observed_gateway_id = _validate_identifier( diff --git a/gateway/hosted_rooms.py b/gateway/hosted_rooms.py index 0d55fe118a..4a823bf0fd 100644 --- a/gateway/hosted_rooms.py +++ b/gateway/hosted_rooms.py @@ -1,9 +1,8 @@ """Durable state for gateway-hosted Bot Mode rooms. -Owns only hosted-room identity and its append-only event log; delivery, relay leasing -and agent turns belong to the relay and the hosted-room driver, so the log composes with -a durable relay without a second transport queue. Callers supply the database path -(production handlers use the gateway's root ``state.db``). +Owns only hosted-room identity and its append-only event log; delivery, relay leasing and agent turns +belong to the relay and the hosted-room driver, so the log composes with a durable relay without a second +transport queue. Callers supply the database path (production handlers use the gateway's root ``state.db``). """ from __future__ import annotations @@ -39,9 +38,8 @@ MAX_DISBANDED_ROOM_TOMBSTONES = 512 DISBANDED_ROOM_RETENTION_SECONDS = 90 * 24 * 60 * 60 MAX_EVENTS_PER_ROOM = 50_000 MAX_ROOM_EVENT_BYTES = 256 * 1024 * 1024 -# Leave substantial headroom below the pre-update state.db snapshot ceiling. -# Event accounting does not include SQLite indexes or repeated room ids, so the -# logical budget must stay well below the physical-file limit. +# Leave substantial headroom below the pre-update state.db snapshot ceiling: event accounting excludes +# SQLite indexes and repeated room ids, so the logical budget must stay well below the physical-file limit. MAX_GATEWAY_EVENT_BYTES = 16 * 1024 * 1024 CONTROL_EVENT_COUNT_RESERVE = 64 CONTROL_EVENT_BYTE_RESERVE = 1024 * 1024 @@ -335,9 +333,8 @@ 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. +# 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 = ( @@ -377,10 +374,9 @@ def _initialize_schema(conn: sqlite3.Connection) -> None: for statement in _SCHEMA_DDL: conn.execute(statement) _migrate_legacy_columns(conn) - # Old schemas kept the final identity tombstone in hosted_rooms itself. - # Copy those identities before bounded history pruning can remove their - # heavier room/event payloads. This compact registry is intentionally - # permanent: a stale coordinate must never name a different Group Chat. + # Old schemas kept the final identity tombstone in hosted_rooms itself. Copy those identities before + # bounded history pruning can remove their heavier room/event payloads. This compact registry is + # intentionally permanent: a stale coordinate must never name a different Group Chat. conn.execute(_RETIRE_FROM_ROOMS.format(where="disbanded_at IS NOT NULL")) _migrate_remote_run_schema(conn) conn.execute("CREATE INDEX IF NOT EXISTS idx_hosted_room_events_cursor ON hosted_room_events(room_id, seq)") @@ -446,8 +442,7 @@ def _room_row(conn: sqlite3.Connection, sql: str, params: tuple[Any, ...], room_ row = conn.execute(sql, params).fetchone() if row is not None: return row - # A retained disband tombstone still has replayable history; the caller simply - # did not opt into reading disbanded rooms. + # A retained disband tombstone still has replayable history; the caller did not opt into disbanded rooms. retained = conn.execute("SELECT 1 FROM hosted_rooms WHERE room_id=?", (room_id,)).fetchone() if retained is None and _is_retired(conn, room_id): raise RoomHistoryExpiredError("Group Chat history expired; room_id remains permanently retired") @@ -918,10 +913,8 @@ def rename_room(db_path: DbPath, *, room_id: Any, event_id: Any, name: Any, now: def append_event( db_path: DbPath, *, room_id: Any, event_id: Any, kind: Any, actor: Any, payload: Any, authority_gateway_id: Any = None, authority_epoch: Any = None, now: float | None = None) -> dict[str, Any]: - """Append one immutable event and allocate its per-room sequence atomically. - - Repeating an ``event_id`` with identical content returns the original; different content fails closed. - """ + """Append one immutable event and allocate its per-room sequence atomically; repeating an ``event_id`` + with identical content returns the original, different content fails closed.""" room_id = _room_id(room_id) event_id = _event_id(event_id) kind = _validate_event_kind(kind) @@ -972,11 +965,9 @@ def _probe(path: Path, table: str, query: str, params: tuple[Any, ...], unavaila def probe_hosted_room(db_path: DbPath, *, room_id: Any) -> bool: - """Check room ownership without creating or migrating the shared store. - - Runs on the synchronous prompt-admission path for older Desktop clients, so it fails fast - under contention instead of blocking the WebSocket reader for SQLite's ten-second timeout. - """ + """Check room ownership without creating or migrating the shared store; runs on the synchronous + prompt-admission path for older Desktop clients, so it fails fast under contention instead of blocking + the WebSocket reader for SQLite's ten-second timeout.""" return _probe( Path(db_path), "hosted_rooms", "SELECT 1 FROM hosted_rooms WHERE room_id=? AND disbanded_at IS NULL LIMIT 1", (_room_id(room_id),), "hosted room ownership is temporarily unavailable") @@ -1019,11 +1010,9 @@ def request_room_stop( def claim_authority( db_path: DbPath, *, room_id: Any, expected_gateway_id: Any, expected_epoch: Any, new_gateway_id: Any, event_id: Any, now: float | None = None) -> dict[str, Any]: - """Fence a verified authority transfer with a compare-and-swap epoch. - - Does not decide *when* takeover is safe: a replicated driver calls it only after its - lease/quorum policy established that the previous owner can no longer commit. - """ + """Fence a verified authority transfer with a compare-and-swap epoch; does not decide *when* takeover is + safe (a replicated driver calls it only after its lease/quorum policy established that the previous + owner can no longer commit).""" room_id = _room_id(room_id) expected_gateway_id = _actor_id(expected_gateway_id, "expected_gateway_id") new_gateway_id = _actor_id(new_gateway_id, "new_gateway_id")