diff --git a/gateway/platforms/api_server_room_dispatch.py b/gateway/platforms/api_server_room_dispatch.py index a7623b2b7c..89d7eea05b 100644 --- a/gateway/platforms/api_server_room_dispatch.py +++ b/gateway/platforms/api_server_room_dispatch.py @@ -28,8 +28,7 @@ async def _ensure_hosted_member_session(self, dispatch: Any) -> str: title = f"Group: {dispatch.room_id}" seed = ( f"{dispatch.home_install_id}\0{dispatch.room_id}\0" - f"{dispatch.member_id}\0{dispatch.target_profile}" - ) + f"{dispatch.member_id}\0{dispatch.target_profile}") session_id = f"room_{hashlib.sha256(seed.encode()).hexdigest()[:32]}" def atomic(conn): @@ -40,18 +39,14 @@ async def _ensure_hosted_member_session(self, dispatch: Any) -> str: return session_id clean_title = db.sanitize_title(title) conflict = conn.execute( - "SELECT id FROM sessions WHERE title=? AND id!=?", (clean_title, session_id) - ).fetchone() + "SELECT id FROM sessions WHERE title=? AND id!=?", (clean_title, session_id)).fetchone() if conflict: raise RuntimeError( "Another group already uses this room title on the target gateway. " - "Rename or migrate that group before retrying." - ) + "Rename or migrate that group before retrying.") conn.execute( "INSERT INTO sessions(id, source, title, hidden, started_at) " - "VALUES(?, 'bot_room', ?, 1, ?)", - (session_id, clean_title, time.time()), - ) + "VALUES(?, 'bot_room', ?, 1, ?)", (session_id, clean_title, time.time())) return session_id return await asyncio.to_thread(db._execute_write, atomic) @@ -70,8 +65,7 @@ def _room_dispatch_error(exc: Exception, *, _openai_error) -> "web.Response": async def _normalize_room_dispatch( - self, request: "web.Request", body: Any, *, _api_server -) -> tuple[Any, "web.Response | None"]: + self, request: "web.Request", body: Any, *, _api_server) -> tuple[Any, "web.Response | None"]: """Validate and normalize a scoped RoomLink dispatch request.""" _openai_error = _api_server._openai_error room_token = self._room_grant_token(request) @@ -80,8 +74,7 @@ async def _normalize_room_dispatch( if not isinstance(body, dict) or set(body) - {"input", "hosted_room_dispatch"}: return body, _json_error( _openai_error, "Room dispatch accepts only input and hosted_room_dispatch.", - code="invalid_room_dispatch", status=400, - ) + code="invalid_room_dispatch", status=400) try: from gateway import hosted_rooms from gateway.hosted_room_peer import GatewayRoomCatalog, HostedMemberDispatch, verify_room_grant diff --git a/gateway/platforms/api_server_room_grants.py b/gateway/platforms/api_server_room_grants.py index 8896430bc3..304961054e 100644 --- a/gateway/platforms/api_server_room_grants.py +++ b/gateway/platforms/api_server_room_grants.py @@ -28,8 +28,7 @@ def _require_unchanged_execution_policy(claims: dict[str, Any], execution_policy def _invalid_room_grant(_openai_error) -> "web.Response": return _json_error( _openai_error, "Room authorization is invalid or expired.", - err_type="gateway_auth_error", code="invalid_room_grant", status=401, - ) + err_type="gateway_auth_error", code="invalid_room_grant", status=401) def _room_grant_error_response(exc: Exception, *, _openai_error) -> "web.Response": @@ -37,8 +36,7 @@ def _room_grant_error_response(exc: Exception, *, _openai_error) -> "web.Respons return _invalid_room_grant(_openai_error) return _json_error( _openai_error, "Room authorization needs to be renewed.", - err_type="gateway_auth_error", code="room_reauthorization_required", status=403, - ) + err_type="gateway_auth_error", code="room_reauthorization_required", status=403) def _hard_expiry(claims: dict[str, Any]) -> float: @@ -52,8 +50,7 @@ def _local_target(claims: dict[str, Any] | None, _api_request_profile) -> tuple[ profile = _api_request_profile.get() or "default" installation_id = hosted_rooms.local_authority_gateway_id() if claims is not None and ( - claims["target_profile"] != profile or claims["target_install_id"] != installation_id - ): + claims["target_profile"] != profile or claims["target_install_id"] != installation_id): raise ValueError("room grant target does not match this profile") return profile, installation_id @@ -68,8 +65,7 @@ def _local_room_catalog(self, profile: str, installation_id: str) -> tuple[dict, catalog = catalog_mapping( installation_id=installation_id, protocol_versions=(PROTOCOL_VERSION,), link_modes=("direct",), persistent_process=True, text=True, attachments=False, target_profile=profile, - execution_policy=execution_policy, - ) + execution_policy=execution_policy) return execution_policy, catalog @@ -78,8 +74,7 @@ def _http_routes(self) -> list[tuple[str, str, Any]]: ("POST", "/v1/room-members/invitations", self._handle_room_member_invitation), ("GET", "/v1/room-members/capabilities", self._handle_room_member_capabilities), ("POST", "/v1/room-members/grants/refresh", self._handle_room_member_grant_refresh), - ("POST", "/v1/room-members/grants/revoke", self._handle_room_member_grant_revoke), - ] + ("POST", "/v1/room-members/grants/revoke", self._handle_room_member_grant_revoke)] def _room_grant_token(request: "web.Request") -> str: @@ -119,8 +114,7 @@ def _room_grant_claims(self, request: "web.Request", *, permission: str) -> dict async def _handle_room_member_invitation( - self, request: "web.Request", *, _openai_error, _api_request_profile -) -> "web.Response": + self, request: "web.Request", *, _openai_error, _api_request_profile) -> "web.Response": """Mint a short-lived room/profile grant for a trusted home gateway.""" auth_err = self._check_auth(request) if auth_err: @@ -133,8 +127,7 @@ async def _handle_room_member_invitation( if set(body) - allowed or not required <= set(body): return _json_error( _openai_error, "Invitation is missing required room authority fields.", - code="invalid_room_invitation", status=400, - ) + code="invalid_room_invitation", status=400) try: from gateway import hosted_rooms from gateway.hosted_room_peer import decode_room_grant, issue_room_grant @@ -150,22 +143,15 @@ async def _handle_room_member_invitation( token = issue_room_grant( self._room_grant_secret(), grant_id=str(body.get("grant_id") or f"grant-{uuid.uuid4().hex}"), - room_id=str(body["room_id"]), - home_install_id=str(body["home_install_id"]), + room_id=str(body["room_id"]), home_install_id=str(body["home_install_id"]), authority_gateway_id=str(body["authority_gateway_id"]), - authority_epoch=int(body["authority_epoch"]), - member_id=str(body["member_id"]), - target_install_id=target_install_id, - target_profile=profile, - execution_policy_digest=execution_policy["policy_digest"], - issued_at=time.time(), - ttl_seconds=ttl, - status_ttl_seconds=status_ttl, - ) + authority_epoch=int(body["authority_epoch"]), member_id=str(body["member_id"]), + target_install_id=target_install_id, target_profile=profile, + execution_policy_digest=execution_policy["policy_digest"], issued_at=time.time(), + ttl_seconds=ttl, status_ttl_seconds=status_ttl) claims = decode_room_grant(self._room_grant_secret(), token, permission="status") hosted_rooms.reserve_peer_room( - hosted_rooms.default_db_path(), claims=claims, expires_at=_hard_expiry(claims) - ) + hosted_rooms.default_db_path(), claims=claims, expires_at=_hard_expiry(claims)) except Exception as exc: return _json_error(_openai_error, str(exc), code="invalid_room_invitation", status=400) return web.json_response({ @@ -179,8 +165,7 @@ async def _handle_room_member_invitation( async def _handle_room_member_capabilities( - self, request: "web.Request", *, _openai_error, _api_request_profile -) -> "web.Response": + self, request: "web.Request", *, _openai_error, _api_request_profile) -> "web.Response": """Verify a scoped grant and return this target's live room catalog.""" try: claims = self._room_grant_claims(request, permission="status") @@ -196,13 +181,11 @@ async def _handle_room_member_capabilities( "authority_epoch": claims["authority_epoch"], "member_id": claims["member_id"], "target_profile": profile, - "catalog": catalog, - }) + "catalog": catalog}) async def _handle_room_member_grant_refresh( - self, request: "web.Request", *, _openai_error, _api_request_profile -) -> "web.Response": + self, request: "web.Request", *, _openai_error, _api_request_profile) -> "web.Response": """Refresh dispatch access without a Desktop or broad gateway key.""" body, error = await self._read_json_body(request) if error: @@ -210,8 +193,7 @@ async def _handle_room_member_grant_refresh( if set(body) - {"ttl_seconds"}: return _json_error( _openai_error, "Grant refresh accepts only ttl_seconds.", - code="invalid_room_grant_refresh", status=400, - ) + code="invalid_room_grant_refresh", status=400) try: from gateway.hosted_room_peer import MAX_DISPATCH_GRANT_TTL_SECONDS, issue_room_grant from gateway.hosted_room_execution_policy import execution_policy_mapping @@ -231,21 +213,14 @@ async def _handle_room_member_grant_refresh( execution_policy = execution_policy_mapping(target_profile=profile) _require_unchanged_execution_policy(claims, execution_policy) token = issue_room_grant( - self._room_grant_secret(), - grant_id=f"grant-refresh-{uuid.uuid4().hex}", - room_id=claims["room_id"], - home_install_id=claims["home_install_id"], + self._room_grant_secret(), grant_id=f"grant-refresh-{uuid.uuid4().hex}", + room_id=claims["room_id"], home_install_id=claims["home_install_id"], authority_gateway_id=claims["authority_gateway_id"], - authority_epoch=int(claims["authority_epoch"]), - member_id=claims["member_id"], - target_install_id=installation_id, - target_profile=profile, + authority_epoch=int(claims["authority_epoch"]), member_id=claims["member_id"], + target_install_id=installation_id, target_profile=profile, execution_policy_digest=execution_policy["policy_digest"], - permissions=claims["permissions"], - issued_at=now, - ttl_seconds=dispatch_ttl, - status_expires_at=hard_expiry, - ) + permissions=claims["permissions"], issued_at=now, ttl_seconds=dispatch_ttl, + status_expires_at=hard_expiry) except Exception as exc: return _room_grant_error_response(exc, _openai_error=_openai_error) return web.json_response({ @@ -253,13 +228,11 @@ async def _handle_room_member_grant_refresh( "grant": token, "expires_at": now + dispatch_ttl, "status_expires_at": hard_expiry, - "execution_policy": execution_policy, - }) + "execution_policy": execution_policy}) async def _handle_room_member_grant_revoke( - self, request: "web.Request", *, _openai_error, _api_request_profile -) -> "web.Response": + self, request: "web.Request", *, _openai_error, _api_request_profile) -> "web.Response": """Revoke exactly the scoped grant authenticating this request.""" body, error = await self._read_json_body(request) if error: @@ -267,8 +240,7 @@ async def _handle_room_member_grant_revoke( if body: return _json_error( _openai_error, "Grant revoke accepts no fields.", - code="invalid_room_grant_revoke", status=400, - ) + code="invalid_room_grant_revoke", status=400) try: from gateway import hosted_rooms @@ -278,8 +250,7 @@ async def _handle_room_member_grant_revoke( claims = _decode_request_grant(self, request, permission="status") _local_target(claims, _api_request_profile) hosted_rooms.revoke_room_grant_scope( - hosted_rooms.default_db_path(), claims=claims, expires_at=_hard_expiry(claims) - ) + hosted_rooms.default_db_path(), claims=claims, expires_at=_hard_expiry(claims)) except Exception: return _invalid_room_grant(_openai_error) return web.json_response({"object": "hermes.room_member.grant.revocation", "revoked": True}) diff --git a/gateway/platforms/api_server_run_idempotency.py b/gateway/platforms/api_server_run_idempotency.py index b1e204fa22..fe437fdf0e 100644 --- a/gateway/platforms/api_server_run_idempotency.py +++ b/gateway/platforms/api_server_run_idempotency.py @@ -18,23 +18,19 @@ TERMINAL_STATUSES = frozenset({"completed", "failed", "cancelled", "interrupted" _SELECT_BY_KEY = ( "SELECT fingerprint, run_id, status_json, owner_pid, owner_started, updated_at " - "FROM run_idempotency WHERE scope=? AND idempotency_key=?" -) + "FROM run_idempotency WHERE scope=? AND idempotency_key=?") _EXTEND_RETENTION_BY_KEY = ( "UPDATE run_idempotency SET retention_until=MAX(retention_until, ?) " - "WHERE scope=? AND idempotency_key=? AND fingerprint=?" -) + "WHERE scope=? AND idempotency_key=? AND fingerprint=?") _EXTEND_RETENTION_BY_RUN = ( "UPDATE run_idempotency SET retention_until=MAX(retention_until, ?) " - "WHERE scope=? AND run_id=?" -) + "WHERE scope=? AND run_id=?") # Columns added after the first schema shipped; applied when missing. _MIGRATIONS = { "owner_pid": "INTEGER NOT NULL DEFAULT 0", "owner_started": "INTEGER NOT NULL DEFAULT 0", "retention_until": "REAL NOT NULL DEFAULT 0", - "acknowledged_at": "REAL", -} + "acknowledged_at": "REAL"} def _encode_status(status: Dict[str, Any]) -> str: @@ -47,8 +43,7 @@ def _record(run_id, status_json, owner_pid, owner_started, updated_at) -> dict[s "status": json.loads(status_json), "owner_pid": int(owner_pid or 0), "owner_started": int(owner_started or 0), - "updated_at": float(updated_at or 0), - } + "updated_at": float(updated_at or 0)} def _outcome(row, fingerprint): @@ -87,9 +82,7 @@ class RunIdempotencyStore: except Exception as exc: logger.warning( "Run idempotency storage is unavailable; falling back to " - "process memory, so replay will not survive a restart: %s", - exc, - ) + "process memory, so replay will not survive a restart: %s", exc) self._conn = sqlite3.connect(":memory:", check_same_thread=False) self._db_path = None from hermes_state import apply_wal_with_fallback @@ -116,8 +109,7 @@ class RunIdempotencyStore: if column not in columns: self._conn.execute(f"ALTER TABLE run_idempotency ADD COLUMN {column} {ddl}") self._conn.execute( - "CREATE UNIQUE INDEX IF NOT EXISTS run_idempotency_run_id ON run_idempotency(run_id)" - ) + "CREATE UNIQUE INDEX IF NOT EXISTS run_idempotency_run_id ON run_idempotency(run_id)") self._conn.commit() self._lock = threading.Lock() self._tighten_permissions() @@ -165,14 +157,11 @@ class RunIdempotencyStore: ") VALUES(?,?,?,?,?,?,?,?,?,?)", ( scope, key, fingerprint, run_id, encoded, - int(owner_pid or 0), int(owner_started or 0), retention_until, now, now, - ), - ) + int(owner_pid or 0), int(owner_started or 0), retention_until, now, now)) self._conn.commit() return "created", { "run_id": run_id, "status": status, "owner_pid": int(owner_pid or 0), - "owner_started": int(owner_started or 0), "updated_at": now, - } + "owner_started": int(owner_started or 0), "updated_at": now} def lookup(self, scope: str, key: str, fingerprint: str, *, retention_until: float = 0): """Return ``missing``, ``reused`` or ``conflict`` without reserving.""" @@ -209,8 +198,7 @@ class RunIdempotencyStore: if terminal: self._conn.execute( "DELETE FROM run_idempotency WHERE scope=? AND idempotency_key=?", - (stale_scope, stale_key), - ) + (stale_scope, stale_key)) def status_for_run(self, scope: str, run_id: str, *, retention_until: float = 0) -> dict[str, Any] | None: """Load one durable run status inside its authenticated scope.""" @@ -221,15 +209,12 @@ class RunIdempotencyStore: self._conn.commit() row = self._conn.execute( "SELECT status_json, owner_pid, owner_started, updated_at " - "FROM run_idempotency WHERE scope=? AND run_id=?", - (scope, run_id), - ).fetchone() + "FROM run_idempotency WHERE scope=? AND run_id=?", (scope, run_id)).fetchone() if row is None: return None return { "status": json.loads(row[0]), "owner_pid": int(row[1] or 0), - "owner_started": int(row[2] or 0), "updated_at": float(row[3] or 0), - } + "owner_started": int(row[2] or 0), "updated_at": float(row[3] or 0)} def extend_retention(self, scope: str, run_id: str, until: float) -> bool: """Persist the latest verified recovery horizon for an active grant.""" @@ -252,8 +237,7 @@ class RunIdempotencyStore: with self._lock: self._conn.execute( "UPDATE run_idempotency SET status_json=?, updated_at=? WHERE run_id=?", - (_encode_status(status), time.time(), run_id), - ) + (_encode_status(status), time.time(), run_id)) self._conn.commit() def close(self) -> None: diff --git a/gateway/platforms/api_server_runs.py b/gateway/platforms/api_server_runs.py index be0a58530a..0a1348fe9d 100644 --- a/gateway/platforms/api_server_runs.py +++ b/gateway/platforms/api_server_runs.py @@ -26,24 +26,21 @@ logger = logging.getLogger("gateway.platforms.api_server") _ROOM_RETENTION_REQUEST_KEY = ( RequestKey("hermes.room_run_retention_until", float) if RequestKey is not None - else "hermes.room_run_retention_until" -) + else "hermes.room_run_retention_until") # Forwarded subagent lifecycle fields; free-text ones are secret-redacted. _SUBAGENT_EVENT_KEYS = ( "goal", "task_count", "task_index", "subagent_id", "child_session_id", "delegation_id", "parent_id", "depth", "model", "tool_count", "status", "summary", "duration_seconds", "input_tokens", "output_tokens", "reasoning_tokens", "api_calls", "cost_usd", "files_read", "files_written", - "output_tail", -) + "output_tail") _SUBAGENT_TEXT_KEYS = ("goal", "summary", "output_tail") # Tool-progress event -> SSE payload fields (tool_name, preview, kwargs); key order is wire format. _FIXED_EVENT_FIELDS = { "tool.started": lambda tool, preview, kw: {"tool": tool, "preview": preview}, "tool.completed": lambda tool, preview, kw: { "tool": tool, "duration": round(kw.get("duration", 0), 3), "error": kw.get("is_error", False)}, - "reasoning.available": lambda tool, preview, kw: {"text": preview or ""}, -} + "reasoning.available": lambda tool, preview, kw: {"text": preview or ""}} def _remember_room_retention(request: "web.Request", claims: dict[str, Any]) -> None: @@ -71,13 +68,6 @@ def _run_not_found(_openai_error, run_id: str) -> "web.Response": return _json_error(_openai_error, f"Run not found: {run_id}", code="run_not_found", status=404) -def _idempotency_conflict(_openai_error) -> "web.Response": - return _json_error( - _openai_error, "Idempotency-Key was already used with a different request payload", - code="idempotency_key_conflict", status=409, - ) - - def _uses_room_run_auth(self, request: "web.Request") -> bool: return request.path.endswith("/v1/runs") and bool(self._room_grant_token(request)) @@ -109,21 +99,18 @@ def _initialize_run_state(self, *, store_factory) -> None: def _http_routes(self) -> list[tuple[str, str, Any]]: return [ - ("POST", "/v1/runs", self._handle_runs), - ("GET", "/v1/runs/{run_id}", self._handle_get_run), + ("POST", "/v1/runs", self._handle_runs), ("GET", "/v1/runs/{run_id}", self._handle_get_run), ("GET", "/v1/runs/{run_id}/events", self._handle_run_events), ("POST", "/v1/runs/{run_id}/approval", self._handle_run_approval), ("POST", "/v1/runs/{run_id}/steer", self._handle_steer_run), - ("POST", "/v1/runs/{run_id}/stop", self._handle_stop_run), - ] + ("POST", "/v1/runs/{run_id}/stop", self._handle_stop_run)] def _idempotency_capabilities(self, *, store_type) -> dict[str, Any]: return { "supported": True, "durable": self._run_idempotency_store.durable, - "retention_seconds": store_type.RETENTION_SECONDS, - } + "retention_seconds": store_type.RETENTION_SECONDS} def _close_run_state(self) -> None: @@ -151,8 +138,7 @@ def _set_run_status(self, run_id: str, status: str, **fields: Any) -> Dict[str, should_persist = ( status != previous_status or status in TERMINAL_STATUSES - or bool(field_names & {"output", "error", "usage", "pending_steer", "session_id"}) - ) + or bool(field_names & {"output", "error", "usage", "pending_steer", "session_id"})) if run_id in self._run_idempotency_ids and should_persist: try: self._run_idempotency_store.update_status(run_id, current) @@ -213,16 +199,13 @@ def _run_idempotency_scope(self, request: "web.Request", *, _api_server) -> str: if self._room_grant_token(request): claims = self._room_grant_claims(request, permission=_room_permission_for(request)) _remember_room_retention(request, claims) - identity = ( - f"{claims['room_id']}\0{claims['home_install_id']}\0" - f"{claims['authority_gateway_id']}\0{claims['authority_epoch']}\0" - f"{claims['member_id']}\0{claims['target_install_id']}\0" - f"{claims['target_profile']}" - ) - return hashlib.sha256(identity.encode()).hexdigest() - profile = _api_server._api_request_profile.get() or "default" - identity = self._expected_api_key() or "unauthenticated-test-listener" - return hashlib.sha256(f"{profile}\0{identity}".encode()).hexdigest() + parts = (claims[k] for k in ( + "room_id", "home_install_id", "authority_gateway_id", "authority_epoch", + "member_id", "target_install_id", "target_profile")) + else: + parts = (_api_server._api_request_profile.get() or "default", + self._expected_api_key() or "unauthenticated-test-listener") + return hashlib.sha256("\0".join(map(str, parts)).encode()).hexdigest() def _check_run_auth(self, request: "web.Request", *, permission: str, _api_server) -> "web.Response | None": @@ -261,21 +244,16 @@ def _durable_run_status(self, request: "web.Request", run_id: str) -> Dict[str, scope = self._run_idempotency_scope(request) record = self._run_idempotency_store.status_for_run( - scope, run_id, retention_until=_room_retention_until(request) - ) + scope, run_id, retention_until=_room_retention_until(request)) if record is None: return None status = dict(record["status"]) if status.get("status") not in TERMINAL_STATUSES and not _owner_alive( - int(record.get("owner_pid") or 0), int(record.get("owner_started") or 0) - ): - status.update({ - "status": "interrupted", - "error": "The gateway restarted before this run settled.", - "last_event": "run.interrupted", - "updated_at": time.time(), - }) + int(record.get("owner_pid") or 0), int(record.get("owner_started") or 0)): + status.update( + status="interrupted", error="The gateway restarted before this run settled.", + last_event="run.interrupted", updated_at=time.time()) self._run_idempotency_store.update_status(run_id, status) self._run_statuses[run_id] = status @@ -305,8 +283,7 @@ def _resolve_conversation_history( if not isinstance(entry, dict) or "role" not in entry or "content" not in entry: return [], instructions, None, _json_error( _openai_error, f"conversation_history[{i}] must have 'role' and 'content' fields", - status=400, - ) + status=400) conversation_history.append({"role": str(entry["role"]), "content": str(entry["content"])}) if previous_response_id: logger.debug("Both conversation_history and previous_response_id provided; using conversation_history") @@ -327,23 +304,29 @@ def _resolve_conversation_history( if isinstance(content, list): # flatten multi-part content blocks to text content = " ".join( part.get("text", "") for part in content - if isinstance(part, dict) and part.get("type") == "text" - ) + if isinstance(part, dict) and part.get("type") == "text") conversation_history.append({"role": msg["role"], "content": str(content)}) return conversation_history, instructions, stored_session_id, None -def _replay_response(self, request: "web.Request", record: dict, gateway_session_key) -> "web.Response": - """202 replay of an already-admitted idempotent run.""" - original_id = str(record["run_id"]) - status = self._durable_run_status(request, original_id) or record["status"] - headers = {"Idempotency-Replayed": "true"} +def _accepted_response(run_id: str, status: str, gateway_session_key, *, replayed: bool) -> "web.Response": + """202 admission response; replays are flagged via ``Idempotency-Replayed``.""" + headers = {"Idempotency-Replayed": "true"} if replayed else {} if gateway_session_key: headers["X-Hermes-Session-Key"] = gateway_session_key return web.json_response( - {"run_id": original_id, "status": status.get("status", "queued"), "replayed": True}, - status=202, headers=headers, - ) + {"run_id": run_id, "status": status, "replayed": replayed}, status=202, headers=headers) + + +def _replay_or_conflict(self, request, outcome, record, gateway_session_key, _openai_error) -> "web.Response": + """409 for a fingerprint conflict, else a 202 replay of the already-admitted run.""" + if outcome == "conflict": + return _json_error( + _openai_error, "Idempotency-Key was already used with a different request payload", + code="idempotency_key_conflict", status=409) + original_id = str(record["run_id"]) + status = self._durable_run_status(request, original_id) or record["status"] + return _accepted_response(original_id, status.get("status", "queued"), gateway_session_key, replayed=True) @dataclass(slots=True) @@ -384,8 +367,7 @@ def _idempotency_key_from(request: "web.Request", _openai_error) -> "tuple[str, if key and (len(key) > 255 or any(ord(ch) < 33 or ord(ch) > 126 for ch in key)): return "", _json_error( _openai_error, "Idempotency-Key must be 1-255 visible ASCII characters", - code="invalid_idempotency_key", status=400, - ) + code="invalid_idempotency_key", status=400) return key, None @@ -408,8 +390,7 @@ def _retire_live_run(self, run_id: str) -> None: """Retire agent/task/approval control state once the executor-backed task is done.""" _forget_run( self, run_id, self._active_run_agents, self._active_run_tasks, - self._run_approval_sessions, self._stopping_run_ids, - ) + self._run_approval_sessions, self._stopping_run_ids) def _drop_run_transport(self, run_id: str) -> None: @@ -458,8 +439,7 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res return _json_error(_openai_error, "No user message found in input", status=400) conversation_history, instructions, stored_session_id, history_err = ( - _resolve_conversation_history(self, body, raw_input, _openai_error=_openai_error) - ) + _resolve_conversation_history(self, body, raw_input, _openai_error=_openai_error)) if history_err is not None: return history_err previous_response_id = body.get("previous_response_id") @@ -470,8 +450,7 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res selection_error = self._request_route_conflict_error( session_id=session_id, gateway_session_key=gateway_session_key, requested_model=agent_overrides.get("requested_model"), - requested_provider=agent_overrides.get("requested_provider"), route=route, - ) + requested_provider=agent_overrides.get("requested_provider"), route=route) if selection_error: return _json_error(_openai_error, selection_error, status=400) @@ -481,12 +460,9 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res if idempotency_key: outcome, record = self._run_idempotency_store.lookup( idempotency_scope, idempotency_key, idempotency_fingerprint, - retention_until=_room_retention_until(request), - ) - if outcome == "conflict": - return _idempotency_conflict(_openai_error) - if outcome == "reused" and record is not None: - return _replay_response(self, request, record, gateway_session_key) + retention_until=_room_retention_until(request)) + if outcome == "conflict" or (outcome == "reused" and record is not None): + return _replay_or_conflict(self, request, outcome, record, gateway_session_key, _openai_error) # Enforce concurrency only for a genuinely new run. limited = self._concurrency_limited_response() @@ -521,16 +497,12 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res outcome, record = self._run_idempotency_store.reserve( idempotency_scope, idempotency_key, idempotency_fingerprint, run_id, initial_status, owner_pid=self._run_owner_pid, owner_started=self._run_owner_started, - retention_until=_room_retention_until(request), - ) + retention_until=_room_retention_until(request)) if outcome != "created": _forget_run( self, run_id, self._run_streams, self._run_streams_created, self._run_approval_sessions, - self._run_statuses, self._run_owners, - ) - if outcome == "conflict": - return _idempotency_conflict(_openai_error) - return _replay_response(self, request, record, gateway_session_key) + self._run_statuses, self._run_owners) + return _replay_or_conflict(self, request, outcome, record, gateway_session_key, _openai_error) self._run_idempotency_ids.add(run_id) launch = _RunLaunch( @@ -545,10 +517,7 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res browser_control_transport_family=_api_server._api_request_browser_control_transport_family.get(), ) _start_run_task(self, launch, _api_server=_api_server) - response_headers = {"X-Hermes-Session-Key": gateway_session_key} if gateway_session_key else {} - return web.json_response( - {"run_id": run_id, "status": "started", "replayed": False}, status=202, headers=response_headers - ) + return _accepted_response(run_id, "started", gateway_session_key, replayed=False) def _start_run_task(self, launch: _RunLaunch, *, _api_server) -> None: @@ -568,14 +537,10 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve """Executor-thread body of one run; returns ``(result, usage)``.""" from gateway.session_context import clear_session_vars from gateway.hosted_room_execution_policy import ( - RoomExecutionPolicy, bind_room_execution_policy, reset_room_execution_policy, - ) + RoomExecutionPolicy, bind_room_execution_policy, reset_room_execution_policy) from tools.approval import ( - register_gateway_notify, - reset_current_session_key, - set_current_session_key, - unregister_gateway_notify, - ) + register_gateway_notify, reset_current_session_key, set_current_session_key, + unregister_gateway_notify) session_id = run.session_id effective_task_id = session_id or run.run_id @@ -593,8 +558,7 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve session_tokens = self._bind_api_server_session( chat_id=session_id or "", session_key=run.approval_session_key, session_id=session_id or "", browser_control_principal=run.browser_control_principal, - browser_control_transport_family=run.browser_control_transport_family, - ) + browser_control_transport_family=run.browser_control_transport_family) if session_tokens: resets.append((session_tokens, clear_session_vars)) if run.room_dispatch is not None: @@ -607,8 +571,7 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve _api_server._publish_turn_process_ownership(agent, effective_task_id) r = agent.run_conversation( user_message=run.user_message, conversation_history=run.conversation_history, - task_id=effective_task_id, - ) + task_id=effective_task_id) finally: # Clear ownership immediately so a later stop/cancel can't reap # background work this run deliberately left running. @@ -617,8 +580,7 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve # bind_declared_conversation), with the same precedence gate. if run.declared_selected: self._bind_declared_conversation( - getattr(agent, "session_id", None) or session_id, run.gateway_session_key - ) + getattr(agent, "session_id", None) or session_id, run.gateway_session_key) try: unregister_gateway_notify(run.approval_session_key) finally: @@ -628,8 +590,7 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve return r, { "input_tokens": getattr(agent, "session_prompt_tokens", 0) or 0, "output_tokens": getattr(agent, "session_completion_tokens", 0) or 0, - "total_tokens": getattr(agent, "session_total_tokens", 0) or 0, - } + "total_tokens": getattr(agent, "session_total_tokens", 0) or 0} def _make_approval_notify(self, run: _RunLaunch, *, _api_server) -> Callable[[Dict[str, Any]], None]: @@ -651,9 +612,7 @@ def _make_approval_notify(self, run: _RunLaunch, *, _api_server) -> Callable[[Di "choices": _api_server._approval_event_choices( smart_denied=bool(event.get("smart_denied")), allow_session=event.get("allow_session") is not False, - allow_permanent=event.get("allow_permanent") is not False, - ), - }) + allow_permanent=event.get("allow_permanent") is not False)}) self._set_run_status(run_id, "waiting_for_approval", last_event="approval.request", approval=event) with suppress(Exception): loop.call_soon_threadsafe(q.put_nowait, event) @@ -696,8 +655,7 @@ async def _execute_run(self, run: _RunLaunch, *, _api_server) -> None: requested_model=run.agent_overrides.get("requested_model"), requested_provider=run.agent_overrides.get("requested_provider"), model_options=run.agent_overrides.get("model_options"), route=run.route, - room_dispatch=run.room_dispatch, room_execution_policy=run.room_execution_policy, - ) + room_dispatch=run.room_dispatch, room_execution_policy=run.room_execution_policy) self._active_run_agents[run_id] = agent approval_notify = _make_approval_notify(self, run, _api_server=_api_server) result, usage = await loop.run_in_executor( @@ -758,9 +716,7 @@ def _release_run_owner_if_forgotten(self, run_id: str) -> None: run_id in table for table in ( self._run_statuses, self._active_run_agents, self._active_run_tasks, - self._run_streams, self._run_approval_sessions, - ) - ): + self._run_streams, self._run_approval_sessions)): return self._run_owners.pop(run_id, None) @@ -804,8 +760,7 @@ def _load_owned_run(self, request: "web.Request", *, _api_server, permission: Op async def _handle_get_run(self, request: "web.Request", *, _api_server) -> "web.Response": """GET /v1/runs/{run_id} — return pollable run status for external UIs.""" _, status, _, _, err = _load_owned_run( - self, request, _api_server=_api_server, permission="status", active_fallback=True - ) + self, request, _api_server=_api_server, permission="status", active_fallback=True) return err or web.json_response(status) @@ -872,8 +827,7 @@ async def _handle_run_approval(self, request: "web.Request", *, _api_server) -> _coerce_request_bool = _api_server._coerce_request_bool _openai_error = _api_server._openai_error run_id, _, _, _, err = _load_owned_run( - self, request, _api_server=_api_server, permission="approve", active_fallback=False - ) + self, request, _api_server=_api_server, permission="approve", active_fallback=False) if err is not None: return err @@ -891,8 +845,7 @@ async def _handle_run_approval(self, request: "web.Request", *, _api_server) -> allowed = {"once", "deny"} if room_scoped else {"once", "session", "always", "deny"} resolve_all = ( _coerce_request_bool(body.get("all"), default=False) - or _coerce_request_bool(body.get("resolve_all"), default=False) - ) + or _coerce_request_bool(body.get("resolve_all"), default=False)) approval_session_key = self._run_approval_sessions.get(run_id) for failed, message, code, status in ( (raw_request_id is not None and (not request_id or len(request_id) > 256), @@ -905,16 +858,14 @@ async def _handle_run_approval(self, request: "web.Request", *, _api_server) -> (room_scoped and not request_id, "Room approvals require the exact request_id.", "approval_request_required", 400), (not approval_session_key, - f"Run has no active approval session: {run_id}", "approval_not_active", 409), - ): + f"Run has no active approval session: {run_id}", "approval_not_active", 409)): if failed: return _json_error(_openai_error, message, code=code, status=status) try: from tools.approval import resolve_gateway_approval resolved = resolve_gateway_approval( - approval_session_key, choice, resolve_all=resolve_all, request_id=request_id or None - ) + approval_session_key, choice, resolve_all=resolve_all, request_id=request_id or None) except Exception as exc: logger.exception("[api_server] approval resolution failed for run %s", run_id) return _json_error(_openai_error, str(exc), status=500) @@ -929,18 +880,15 @@ async def _handle_run_approval(self, request: "web.Request", *, _api_server) -> return web.json_response({ "object": "hermes.run.approval_response", "run_id": run_id, - "choice": choice, - **request_id_field, - "resolved": resolved, - }) + "choice": choice, **request_id_field, + "resolved": resolved}) async def _handle_steer_run(self, request: "web.Request", *, _api_server) -> "web.Response": """POST /v1/runs/{run_id}/steer — inject guidance into a running agent.""" _openai_error = _api_server._openai_error run_id, status, agent, _, err = _load_owned_run( - self, request, _api_server=_api_server, permission=None, active_fallback=False - ) + self, request, _api_server=_api_server, permission=None, active_fallback=False) if err is not None: return err # Only genuinely running runs are steerable. /stop retains agent/task refs @@ -949,8 +897,7 @@ async def _handle_steer_run(self, request: "web.Request", *, _api_server) -> "we if status.get("status") != "running" or not hasattr(agent, "steer"): return _json_error( _openai_error, f"Run is not currently accepting steer input: {run_id}", - code="run_not_accepting_steer", status=409, - ) + code="run_not_accepting_steer", status=409) body, err = await self._read_json_body(request) if err: @@ -960,8 +907,7 @@ async def _handle_steer_run(self, request: "web.Request", *, _api_server) -> "we if not steer_text: return _json_error( _openai_error, "Missing non-empty steer text; expected 'input', 'message', or 'text'.", - code="invalid_steer_input", status=400, - ) + code="invalid_steer_input", status=400) try: accepted = bool(agent.steer(steer_text)) @@ -980,8 +926,7 @@ async def _handle_stop_run(self, request: "web.Request", *, _api_server) -> "web """POST /v1/runs/{run_id}/stop — interrupt a running agent.""" _openai_error = _api_server._openai_error run_id, status, agent, task, err = _load_owned_run( - self, request, _api_server=_api_server, permission="stop", active_fallback=True - ) + self, request, _api_server=_api_server, permission="stop", active_fallback=True) if err is not None: return err if status.get("status") in TERMINAL_STATUSES: @@ -990,8 +935,7 @@ async def _handle_stop_run(self, request: "web.Request", *, _api_server) -> "web if agent is None and task is None: return _json_error( _openai_error, f"Run is not active in this gateway process: {run_id}", - code="run_not_active", status=409, - ) + code="run_not_active", status=409) self._set_run_status(run_id, "stopping", last_event="run.stopping") self._stopping_run_ids.add(run_id) @@ -1021,8 +965,7 @@ def _sweep_orphaned_runs_once(self, now: Optional[float] = None) -> None: stale = [ run_id for run_id, created_at in list(self._run_streams_created.items()) - if now - created_at > self._RUN_STREAM_TTL and run_id not in self._run_stream_subscribers - ] + if now - created_at > self._RUN_STREAM_TTL and run_id not in self._run_stream_subscribers] for run_id in stale: logger.debug("[api_server] sweeping expired run transport %s", run_id) task = self._active_run_tasks.get(run_id) @@ -1037,7 +980,6 @@ def _sweep_orphaned_runs_once(self, now: Optional[float] = None) -> None: run_id for run_id, status in list(self._run_statuses.items()) if status.get("status") in {"completed", "failed", "cancelled"} - and now - float(status.get("updated_at", 0) or 0) > self._RUN_STATUS_TTL - ] + and now - float(status.get("updated_at", 0) or 0) > self._RUN_STATUS_TTL] for run_id in stale_statuses: _forget_run(self, run_id, self._run_statuses, self._run_idempotency_ids)