Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
519 changes: 170 additions & 349 deletions .agents/skills/custom-codereview-guide.md

Large diffs are not rendered by default.

24 changes: 24 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,30 @@ The usual flow is SDK/Agent Server → OpenAPI contract → `clients/typescript`

All pull requests must comply with [`.agents/skills/custom-codereview-guide.md`](.agents/skills/custom-codereview-guide.md), in addition to the repository's contribution requirements and CI checks.

## Review-Facing Implementation Checklist

Code should satisfy the repository review checkpoints before the PR is opened:

- Trace cross-layer changes through every affected public entry point, including
factories, constructors, registries, serialization, REST/WebSocket transport,
`clients/typescript/`, and create/resume/fork paths. Do not add a field or
option at one layer while another supported path drops or ignores it.
- Treat public Python and server APIs, defaults, serialized events, persisted
settings, and stored conversations as compatibility surfaces. Use the
deprecation, schema-version migration, and golden-fixture mechanisms described
below instead of one-off shims.
- Give tasks, processes, connections, plugins, event loops, and persistent
artifacts an owner. Cancellation must stop underlying work; close and rollback
paths must cover success, failure, and cancellation; shared conversation state
must use its existing synchronization mechanism.
- Route credentials through `openhands.sdk.utils.pydantic_secrets` and verify the
full input, serialization, persistence, logging, resume, and delivery path.
Never introduce a parallel redaction or secret-sentinel implementation.
- Verify imports, executables, dependency installation, paths, and process
cleanup in every affected production artifact, including the packaged Agent
Server, Docker images, and relevant host platforms. A mocked unit test alone
does not validate a packaging or installation change.

## Repository Memory
- Async LLM completions propagate through the full call chain: `LLM.acompletion()`/`LLM.aresponses()` → `_atransport_call()` (litellm `acompletion`/`aresponses`) → `RetryMixin.retry_decorator()` (tenacity `retry`, which wraps coroutines natively — there is no separate async retry path) → condenser `acondense()` → `Agent.astep()` → `LocalConversation.arun()` → `EventService.run()`. Every async method has a sync counterpart; base classes provide default delegations to sync so custom subclasses work without changes. Token callbacks use `AnyTokenCallbackType` (union of sync/async) with `_invoke_token_callback()` for transparent dispatch.
- `conversation.interrupt()` cancels in-flight `arun()` by cancelling the tracked `_arun_task`. `asyncio.CancelledError` propagates through all layers (LLM HTTP stream → agent step → conversation loop) without needing per-layer interrupt APIs, because LLM and Agent are frozen/stateless Pydantic models that may be shared across conversations. `arun()` catches `CancelledError`, sets status to `PAUSED`, and emits `InterruptEvent`. The agent-server exposes this via `EventService.interrupt()` → `ConversationService.interrupt_conversation()` → `POST /{conversation_id}/interrupt`.
Expand Down
68 changes: 59 additions & 9 deletions openhands-agent-server/openhands/agent_server/mcp_oauth_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

import asyncio
import copy
import threading
from collections.abc import Mapping, Sequence
from dataclasses import dataclass
from typing import Any, NamedTuple, SupportsFloat
Expand All @@ -29,6 +30,7 @@
ToolsReconciledCallback,
create_mcp_tools,
)
from openhands.sdk.utils.cipher import Cipher


logger = get_logger(__name__)
Expand Down Expand Up @@ -83,7 +85,33 @@ def _state_field_for_fastmcp_key(


class MCPSettingsOAuthTokenStore:
"""FastMCP OAuth token storage persisted inside settings MCP servers."""
"""FastMCP OAuth token storage persisted inside settings MCP servers.

``seed_mcp_config`` covers OAuth servers that are not in this server's
settings store, such as the ones a hosted deployment passes inline on the
agent (their settings live on the remote app server): the ``auth.state``
they carry is served to FastMCP, and tokens refreshed during the
conversation are kept in memory for the lifetime of the store instead of
being dropped. Servers found in settings always take precedence.
"""

def __init__(
self,
*,
seed_mcp_config: Mapping[str, MCPServer] | None = None,
cipher: Cipher | None = None,
):
self._seeded: dict[str, MCPOAuthState] = {}
self._seeded_lock = threading.Lock()
for server in (seed_mcp_config or {}).values():
if server.url is None or server.oauth_auth is None:
continue
state = server.initial_oauth_state(cipher=cipher) or MCPOAuthState()
self._seeded[server.url.rstrip("/")] = state

def _seeded_state(self, key: str) -> MCPOAuthState | None:
with self._seeded_lock:
return self._seeded.get(_server_url_from_fastmcp_key(key))

def _get_entry_sync(
self, key: str, collection: str | None
Expand All @@ -94,13 +122,13 @@ def _get_entry_sync(

store = get_settings_store()
settings = store.load()
if settings is None:
return None, None

mcp_config = settings.agent_settings.mcp_config
mcp_config = settings.agent_settings.mcp_config if settings else {}
match = _find_matching_oauth_server(mcp_config, key)
if match is None:
return None, None
seeded = self._seeded_state(key)
if seeded is None:
return None, None
return seeded.get_token_storage_value(field), None
_, _, auth = match
return (auth.state or MCPOAuthState()).get_token_storage_value(field), None

Expand Down Expand Up @@ -133,6 +161,14 @@ def apply_update(settings: PersistedSettings) -> PersistedSettings:
mcp_config = settings.agent_settings.mcp_config
match = _find_matching_oauth_server(mcp_config, key)
if match is None:
server_url = _server_url_from_fastmcp_key(key)
with self._seeded_lock:
seeded = self._seeded.get(server_url)
if seeded is not None:
self._seeded[server_url] = seeded.with_token_storage_value(
field, stored_value
)
return settings
logger.warning(
"Could not persist MCP OAuth state: no configured MCP "
"server matches FastMCP key %r",
Expand Down Expand Up @@ -188,6 +224,12 @@ def apply_update(settings: PersistedSettings) -> PersistedSettings:
mcp_config = settings.agent_settings.mcp_config
match = _find_matching_oauth_server(mcp_config, key)
if match is None:
server_url = _server_url_from_fastmcp_key(key)
with self._seeded_lock:
seeded = self._seeded.get(server_url)
if seeded is not None:
seeded, deleted = seeded.without_token_storage_value(field)
self._seeded[server_url] = seeded
return settings
server_name, server, auth = match
state, deleted = (
Expand Down Expand Up @@ -330,7 +372,13 @@ async def delete_many(

@dataclass(frozen=True, slots=True)
class SettingsBackedMCPToolProvider:
"""Create MCP tools with FastMCP OAuth state persisted in settings."""
"""Create MCP tools with FastMCP OAuth state persisted in settings.

OAuth servers absent from settings (passed inline on the agent) fall back
to the OAuth state they carry; see ``MCPSettingsOAuthTokenStore``.
"""

cipher: Cipher | None = None

def create_tools(
self,
Expand All @@ -343,7 +391,9 @@ def create_tools(
return create_mcp_tools(
mcp_config,
timeout,
mcp_oauth_token_storage=MCPSettingsOAuthTokenStore(),
mcp_oauth_token_storage=MCPSettingsOAuthTokenStore(
seed_mcp_config=mcp_config, cipher=self.cipher
),
on_tools_changed=on_tools_changed,
on_tools_reconciled=on_tools_reconciled,
)
Expand All @@ -360,4 +410,4 @@ def create_settings_backed_mcp_tool_provider(
"(no OH_SECRET_KEY configured). Configure OH_SECRET_KEY for "
"production deployments."
)
return SettingsBackedMCPToolProvider()
return SettingsBackedMCPToolProvider(cipher=config.cipher)
8 changes: 8 additions & 0 deletions openhands-sdk/openhands/sdk/context/agent_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -260,6 +260,14 @@ def _validate_skills(cls, v: list[Skill], _info):
@model_validator(mode="after")
def _load_auto_skills(self):
"""Load user/legacy-public skills if enabled, then apply ``disabled_skills``."""
return self.resolve_auto_skills()

def resolve_auto_skills(self) -> AgentContext:
"""Resolve ``load_*_skills`` into ``skills`` and apply ``disabled_skills``.

Exposed because ``model_copy`` skips validators: callers that rebuild a
context that way must re-run this or the deny-list silently won't apply.
"""
include_public = self.load_public_skills
if self.load_user_skills or include_public:
auto_skills = load_available_skills(
Expand Down
5 changes: 4 additions & 1 deletion openhands-sdk/openhands/sdk/llm/llm.py
Original file line number Diff line number Diff line change
Expand Up @@ -1129,8 +1129,11 @@ def _process_stream_event(
):
delta = event.delta
if delta:
# ModelResponseStream mints a fresh id per instance, and a
# changed chunk id reads as a retry (StreamContext._emit_delta).
delta_chunk = ModelResponseStream(
choices=[StreamingChoices(delta=Delta(content=delta))]
id=event.item_id,
choices=[StreamingChoices(delta=Delta(content=delta))],
)

return output_item, delta_chunk
Expand Down
23 changes: 19 additions & 4 deletions openhands-sdk/openhands/sdk/workspace/remote/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -886,6 +886,7 @@ def load_skills_from_agent_server(
load_project: bool = True,
load_org: bool = True,
timeout: float = 60.0,
base_context: "AgentContext | None" = None,
) -> tuple[list["Skill"], "AgentContext"]:
"""Load skills via the agent-server's /api/skills endpoint.

Expand All @@ -906,6 +907,11 @@ def load_skills_from_agent_server(
load_project: Load project skills from workspace directories.
load_org: Load organization-level skills.
timeout: Request timeout in seconds.
base_context: Existing AgentContext to preserve. All of its
fields survive except `skills` and `load_public_skills`,
which this method always sets based on whether skills
were found. Defaults to None, which starts from a fresh
AgentContext — today's behavior.

Returns:
Tuple of (list of Skill objects, AgentContext).
Expand Down Expand Up @@ -958,11 +964,20 @@ def load_skills_from_agent_server(
if loaded_skills:
logger.debug(f"Skills: {[s.name for s in loaded_skills]}")

# Create AgentContext - fall back to public skills if none loaded
# Update `base_context` (or start fresh if none given) with the
# newly loaded skills — every other field the caller configured is
# preserved. Fall back to public skills if none loaded.
base = base_context if base_context is not None else AgentContext()
if loaded_skills:
agent_context = AgentContext(skills=loaded_skills, load_public_skills=False)
agent_context = base.model_copy(
update={"skills": loaded_skills, "load_public_skills": False}
)
else:
logger.warning("No skills loaded, falling back to public skills")
agent_context = AgentContext(skills=[], load_public_skills=True)
agent_context = base.model_copy(
update={"skills": [], "load_public_skills": True}
)

return loaded_skills, agent_context
# ``model_copy`` skips validators, so re-run the resolution that
# applies ``load_*_skills`` and the ``disabled_skills`` deny-list.
return loaded_skills, agent_context.resolve_auto_skills()
2 changes: 2 additions & 0 deletions openhands-workspace/openhands/workspace/cloud/workspace.py
Original file line number Diff line number Diff line change
Expand Up @@ -939,6 +939,7 @@ def load_skills_from_agent_server(
load_project: bool = True,
load_org: bool = True,
timeout: float = 60.0,
base_context: AgentContext | None = None,
) -> tuple[list[Skill], AgentContext]:
"""Load skills from the agent server.

Expand All @@ -951,6 +952,7 @@ def load_skills_from_agent_server(
load_project=load_project,
load_org=load_org,
timeout=timeout,
base_context=base_context,
)

def _call_skills_api(
Expand Down
123 changes: 123 additions & 0 deletions tests/agent_server/test_mcp_oauth_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -385,3 +385,126 @@ def callback(client, tools):
provider.create_tools(config, on_tools_reconciled=callback)

assert mock_create.call_args.kwargs["on_tools_reconciled"] is callback


def _inline_oauth_server(access_token: str):
"""One OAuth server carrying its state inline, as a hosted agent passes it."""
return coerce_mcp_config(
{
"mail": {
"url": "https://mcp.example.com/mcp",
"auth": {
"strategy": "oauth2",
"authentication": {"type": "oauth", "client_auth_method": "none"},
"state": {"tokens": {"access_token": access_token}},
},
}
}
)


@pytest.mark.asyncio
async def test_mcp_oauth_token_store_serves_inline_state_for_servers_absent_from_settings( # noqa: E501
tmp_path: Path,
):
reset_stores()
try:
# Arrange: the sandbox settings store knows no MCP servers at all.
config = Config(
session_api_keys=[],
conversations_path=tmp_path / "conversations",
secret_key=SecretStr("mcp-oauth-test-key"),
)
settings_store = get_settings_store(config)
settings_store.save(PersistedSettings())
store = MCPSettingsOAuthTokenStore(
seed_mcp_config=_inline_oauth_server("agent-access-token"),
cipher=config.cipher,
)
key = "https://mcp.example.com/mcp/tokens"

# Act
initial = await store.get(key=key, collection="mcp-oauth-token")
await store.put(
key=key,
value={"access_token": "refreshed-access-token"},
collection="mcp-oauth-token",
)
refreshed = await store.get(key=key, collection="mcp-oauth-token")
deleted = await store.delete(key=key, collection="mcp-oauth-token")

# Assert: served from the inline state, refreshed in memory, never
# written into settings.
assert initial == {"access_token": "agent-access-token"}
assert refreshed == {"access_token": "refreshed-access-token"}
assert deleted is True
assert await store.get(key=key, collection="mcp-oauth-token") is None
loaded = settings_store.load()
assert loaded is not None
assert not loaded.agent_settings.mcp_config
finally:
reset_stores()


@pytest.mark.asyncio
async def test_mcp_oauth_token_store_prefers_settings_over_inline_state(
tmp_path: Path,
):
reset_stores()
try:
# Arrange: the same server exists in settings with its own tokens.
config = Config(
session_api_keys=[],
conversations_path=tmp_path / "conversations",
secret_key=SecretStr("mcp-oauth-test-key"),
)
settings = PersistedSettings()
settings.agent_settings = settings.agent_settings.model_copy(
update={"mcp_config": _inline_oauth_server("settings-access-token")}
)
get_settings_store(config).save(settings)
store = MCPSettingsOAuthTokenStore(
seed_mcp_config=_inline_oauth_server("agent-access-token"),
cipher=config.cipher,
)

# Act
value = await store.get(
key="https://mcp.example.com/mcp/tokens", collection="mcp-oauth-token"
)

# Assert
assert value == {"access_token": "settings-access-token"}
finally:
reset_stores()


@pytest.mark.asyncio
async def test_settings_backed_provider_seeds_token_store_from_agent_mcp_config(
tmp_path: Path,
):
reset_stores()
try:
# Arrange
config = Config(
session_api_keys=[],
conversations_path=tmp_path / "conversations",
secret_key=SecretStr("mcp-oauth-test-key"),
)
get_settings_store(config).save(PersistedSettings())
provider = create_settings_backed_mcp_tool_provider(config)
mcp_config = _inline_oauth_server("agent-access-token")

# Act
with patch(
"openhands.agent_server.mcp_oauth_store.create_mcp_tools"
) as mock_create:
provider.create_tools(mcp_config)

# Assert: the store handed to FastMCP already knows the agent's tokens.
storage = mock_create.call_args.kwargs["mcp_oauth_token_storage"]
assert await storage.get(
key="https://mcp.example.com/mcp/tokens", collection="mcp-oauth-token"
) == {"access_token": "agent-access-token"}
finally:
reset_stores()
Loading
Loading