refactor(gateway/platforms): runs — unify 202 admission/replay responses, scope hashing, AST-neutral layout packing across the four /v1/runs + room modules

This commit is contained in:
Teknium
2026-09-02 19:12:53 -07:00
parent e282b09973
commit a12a72314d
4 changed files with 117 additions and 227 deletions
+6 -13
View File
@@ -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
+27 -56
View File
@@ -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})
+13 -29
View File
@@ -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:
+71 -129
View File
@@ -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)