refactor(gateway): reflow hosted-room module/function docstrings to <=120 cols (every sentence kept)
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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 <set> WHERE room_id=? AND task_id=? AND <fence>";
|
||||
# the fence names the expected status plus the generations that must not have moved.
|
||||
# Every transition is "UPDATE ... SET <set> WHERE room_id=? AND task_id=? AND <fence>"; 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 (
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
+19
-30
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user