diff --git a/backend/druks/agents.py b/backend/druks/agents.py index eb5f7dda..8cc745c3 100644 --- a/backend/druks/agents.py +++ b/backend/druks/agents.py @@ -338,9 +338,9 @@ async def _run( # same servers. subject = await workflow.subject workspace_class = workflow.workspace_class - mcp_servers, mcp_refs = (), [] + mcp_servers, mcp_secret_refs = (), [] if self.include_mcp: - mcp_servers, mcp_refs = await workspace_class.get_all_mcp_servers( + mcp_servers, mcp_secret_refs = await workspace_class.get_all_mcp_servers( session, subject, workflow.account_id ) refs = [ @@ -354,7 +354,7 @@ async def _run( ) for secret in await workspace_class.get_secrets(subject) ), - *mcp_refs, + *mcp_secret_refs, ] host_id = await workflow._lease_host(session, config, refs) diff --git a/backend/druks/chat/bridge.mjs b/backend/druks/chat/bridge.mjs index 64befb24..1c7440c7 100644 --- a/backend/druks/chat/bridge.mjs +++ b/backend/druks/chat/bridge.mjs @@ -160,7 +160,15 @@ class Conversation { if (!bearer) throw new Error("The Druks MCP placeholder is missing."); const setup = { cwd: this.root, - mcpServers: [{ name: "druks", type: "http", url: request.mcpUrl, headers: [{ name: "Authorization", value: bearer }, ...request.headers] }], + mcpServers: [ + { name: "druks", type: "http", url: request.mcpUrl, headers: [{ name: "Authorization", value: bearer }, ...request.headers] }, + ...request.mcpServers.map(server => ({ + name: server.name, + type: "http", + url: server.url, + headers: Object.entries(server.headers).map(([name, value]) => ({ name, value: fillPlaceholders(value) })), + })), + ], _meta: request.meta, }; const session = this.state.sessionId diff --git a/backend/druks/chat/service.py b/backend/druks/chat/service.py index e39799d1..2cd7ecff 100644 --- a/backend/druks/chat/service.py +++ b/backend/druks/chat/service.py @@ -33,6 +33,7 @@ from druks.redis import get_client from druks.sandbox.client import sandbox_client from druks.sandbox.constants import SANDBOX_HOST_LEASE_SECONDS +from druks.sandbox.datastructures import McpServer from druks.sandbox.exceptions import HostGone from druks.sandbox.host import Host from druks.sandbox.layout import get_remote_home, get_work_root @@ -103,8 +104,10 @@ async def get_agent( async def get_sandbox( session: AsyncSession, account_id: str, + *, config: AgentConfig, allowed_tools: AllowedTools, + mcp_secret_refs: list[SecretRef], ) -> tuple[Host, SandboxIdentity]: """The account's sandbox for a new turn: a live sandbox that holds the Chat agent's current secrets, or a new one.""" @@ -112,6 +115,7 @@ async def get_sandbox( token = await get_druks_account_token(session, account_id, allowed_tools, name=CHAT_KEY_NAME) refs = [ *config.secret_refs, + *mcp_secret_refs, SecretRef( name=get_bearer_token_env_var(server.name).lower(), secret_id=token.id, @@ -204,11 +208,31 @@ async def deliver_pending(session: AsyncSession, conversation: Conversation) -> f"Chat runs on {adapters}. The Bot's settings select " f"{config.harness_class.name}. Set its harness to one of them." ) - host, identity = await get_sandbox(session, conversation.account_id, config, tools) + # A bot serves outside people, so only the operator reaches the enabled MCP servers. + mcp_servers, mcp_secret_refs = (), [] + if conversation.account.kind == AccountKind.OPERATOR: + # A server the operator has not connected must not stop the chat. + mcp_servers, mcp_secret_refs = await Workspace.get_all_mcp_servers( + session, None, conversation.account_id, skip_unauthenticated=True + ) + host, identity = await get_sandbox( + session, + conversation.account_id, + config=config, + allowed_tools=tools, + mcp_secret_refs=mcp_secret_refs, + ) try: bridge = Bridge(host) turn = await send_turn( - session, conversation, message, bridge, identity, config, prompt + session, + conversation, + message, + bridge=bridge, + identity=identity, + config=config, + prompt=prompt, + mcp_servers=mcp_servers, ) if turn: await follow_turn(session, conversation, turn, bridge) @@ -234,10 +258,12 @@ async def send_turn( session: AsyncSession, conversation: Conversation, message: Message, + *, bridge: Bridge, identity: SandboxIdentity, config: AgentConfig, prompt: str, + mcp_servers: tuple[McpServer, ...], ) -> Message | None: """Start the agent, transcribe the pending voice notes, and send the pending messages. Return the turn's message, unless a Stop or pause came first.""" @@ -271,6 +297,14 @@ async def send_turn( mcpUrl=server.url, bearerVariable=get_bearer_token_env_var(server.name), headers=headers, + mcpServers=[ + { + "name": mcp_server.name, + "url": mcp_server.url, + "headers": mcp_server.get_request_headers(), + } + for mcp_server in mcp_servers + ], **config.harness_class.get_acp_session( account_type, config.model, prompt, config.identity, sandbox_home, conversation_root ), diff --git a/backend/druks/harnesses/claude.py b/backend/druks/harnesses/claude.py index 22595649..95d6043a 100644 --- a/backend/druks/harnesses/claude.py +++ b/backend/druks/harnesses/claude.py @@ -211,20 +211,16 @@ def _mcp_flags(self, servers: tuple[McpServer, ...]) -> tuple[str, ...]: Claude expands ``${VAR}`` refs in header values from the run env at connect time, so no secret — the bearer or a secret header — ever lands in the emitted config; only non-secret values ride inline.""" - if not servers: - return () - entries = {} - for server in servers: - headers = dict(server.headers) - if server.bearer_token_env_var: - headers["Authorization"] = f"Bearer ${{{server.bearer_token_env_var}}}" - for header, env_var in server.env_headers.items(): - headers[header] = f"${{{env_var}}}" - entry: dict[str, object] = {"type": "http", "url": server.url} - if headers: - entry["headers"] = headers - entries[server.name] = entry - return ("--mcp-config", json.dumps({"mcpServers": entries})) + if servers: + entries = {} + for server in servers: + headers = server.get_request_headers() + entry: dict[str, object] = {"type": "http", "url": server.url} + if headers: + entry["headers"] = headers + entries[server.name] = entry + return ("--mcp-config", json.dumps({"mcpServers": entries})) + return () @classmethod def get_secret_refs(cls, subscription: VaultSecret) -> list[SecretRef]: diff --git a/backend/druks/harnesses/pi.py b/backend/druks/harnesses/pi.py index 389b2df4..044d3efc 100644 --- a/backend/druks/harnesses/pi.py +++ b/backend/druks/harnesses/pi.py @@ -74,11 +74,7 @@ async def build_invocation( extension_body = _DRUKS_OUTPUT_TEMPLATE.read_text() mcp = {} for server in mcp_servers: - headers = dict(server.headers) - if server.bearer_token_env_var: - headers["Authorization"] = f"Bearer ${{{server.bearer_token_env_var}}}" - for header, env_var in server.env_headers.items(): - headers[header] = f"${{{env_var}}}" + headers = server.get_request_headers() # Druks owns MCP authentication, so the adapter must not start headless OAuth. entry: dict[str, object] = {"url": server.url, "auth": False} if headers: diff --git a/backend/druks/sandbox/datastructures.py b/backend/druks/sandbox/datastructures.py index 52b8ebf6..710348e0 100644 --- a/backend/druks/sandbox/datastructures.py +++ b/backend/druks/sandbox/datastructures.py @@ -147,6 +147,15 @@ class McpServer: # Secret declared headers: header name -> the env var carrying its value. env_headers: dict[str, str] = field(default_factory=dict) + def get_request_headers(self) -> dict[str, str]: + """The headers a request to the server carries, each secret one as a ${VAR} placeholder.""" + headers = dict(self.headers) + if self.bearer_token_env_var: + headers["Authorization"] = f"Bearer ${{{self.bearer_token_env_var}}}" + for header, env_var in self.env_headers.items(): + headers[header] = f"${{{env_var}}}" + return headers + @dataclass(frozen=True) class SandboxSecret: diff --git a/backend/druks/workspaces.py b/backend/druks/workspaces.py index 877c0a41..3afeec72 100644 --- a/backend/druks/workspaces.py +++ b/backend/druks/workspaces.py @@ -166,13 +166,19 @@ async def run_agent(self, *, account_id: str | None, **kwargs: Any) -> AgentResu @classmethod async def get_all_mcp_servers( - cls, session: AsyncSession, subject: Any, account_id: str | None + cls, + session: AsyncSession, + subject: Any, + account_id: str | None, + *, + skip_unauthenticated: bool = False, ) -> tuple[tuple[McpServer, ...], list[SecretRef]]: """The MCP servers a box of this workspace reaches, as the harness names them, and the secret refs for the box's entries, one per bearer and per secret header. The workspace's servers come first and own their names: - a same-named registry entry is neither resolved nor delivered. A server - that cannot authenticate fails here, before the box.""" + a same-named registry entry is neither resolved nor delivered. A registry + server that cannot authenticate fails here, before the box, unless + ``skip_unauthenticated`` leaves it out.""" workspace_servers = await cls.get_mcp_servers(subject) workspace_names = {server.name for server in workspace_servers} if len(workspace_names) != len(workspace_servers): @@ -204,7 +210,11 @@ async def get_all_mcp_servers( ) owner = await Account.get_secrets_owner(session, account_id) run_account = owner.id if owner else None - for server in await mcp_models.McpServer.list_enabled(session): + registry = await mcp_models.McpServer.list_enabled(session) + if skip_unauthenticated: + resolved = await mcp_models.McpServer.get_resolved(session, run_account) + registry = [server for server in registry if resolved[server["name"]]["has_token"]] + for server in registry: name = server["name"] if name in workspace_names: continue diff --git a/backend/tests/chat/bridge.test.mjs b/backend/tests/chat/bridge.test.mjs index ab593031..980e52fe 100644 --- a/backend/tests/chat/bridge.test.mjs +++ b/backend/tests/chat/bridge.test.mjs @@ -61,7 +61,7 @@ test("The bridge streams detached turns, isolates archives, cancels, and reloads "newSession: async p => { cwd=p.cwd; sessionId=randomUUID();", " if (p._meta.harness !== 'options') throw Error('meta');", " if (p.mcpServers[0].headers[0].value !== 'placeholder') throw Error('token');", - " fs.writeFileSync(path.join(cwd, 'setup.json'), JSON.stringify(p.mcpServers[0].headers));", + " fs.writeFileSync(path.join(cwd, 'setup.json'), JSON.stringify(p.mcpServers));", " return {sessionId, configOptions: [{id: 'model'}, {id: 'effort'}, {id: 'fast'}]}; },", "loadSession: async p => { cwd=p.cwd; sessionId=p.sessionId; memory=fs.readFileSync(transcript(), 'utf8');", " await client.sessionUpdate({sessionId,update:{sessionUpdate:'agent_message_chunk',content:{type:'text',text:'replay'}}}); return {configOptions: [{id: 'model'}, {id: 'effort'}, {id: 'fast'}]}; },", @@ -115,16 +115,19 @@ test("The bridge streams detached turns, isolates archives, cancels, and reloads mode: "bypassPermissions", model: "claude-opus-4-7", meta: { harness: "options" }, env: {}, files: {}, options: { model: "model", effort: "effort", fast: "fast" }, effort: "high", fastMode: true, bearerVariable: "MCP_DRUKS_TOKEN", mcpUrl: "https://hooks.example.com/mcp", - headers: [], + headers: [], mcpServers: [], }); const status = conversation => request(port, { method: "status", conversationId: conversation }); const first = await launch(home); assert.equal((await request(port, start(id(1)))).ok, true); const conversationHeader = { name: "X-Druks-Conversation", value: id(2) }; const login = path.join(home, ".fake/login.json"); - assert.equal((await request(port, { ...start(id(2)), headers: [conversationHeader], files: { [login]: '{"token": "${MCP_DRUKS_TOKEN}"}' } })).ok, true); - const headers = JSON.parse(await fs.readFile(path.join(home, "work", "chat", id(2), "setup.json"), "utf8")); - assert.deepEqual(headers[1], conversationHeader); + const linear = { name: "linear", url: "https://mcp.linear.app/mcp", headers: { Authorization: "Bearer ${MCP_DRUKS_TOKEN}" } }; + assert.equal((await request(port, { ...start(id(2)), headers: [conversationHeader], mcpServers: [linear], files: { [login]: '{"token": "${MCP_DRUKS_TOKEN}"}' } })).ok, true); + const servers = JSON.parse(await fs.readFile(path.join(home, "work", "chat", id(2), "setup.json"), "utf8")); + assert.deepEqual(servers[0].headers[1], conversationHeader); + // A registry server's header names its placeholder, and the bridge fills it the same way. + assert.deepEqual(servers[1], { name: "linear", type: "http", url: linear.url, headers: [{ name: "Authorization", value: "Bearer placeholder" }] }); // A file names a placeholder variable, and the bridge fills it from the sandbox environment. assert.equal(await fs.readFile(login, "utf8"), '{"token": "placeholder"}'); assert.equal((await request(port, { ...start(id(1)), model: "claude-sonnet-5", effort: "low", fastMode: false })).ok, true); diff --git a/backend/tests/test_chat.py b/backend/tests/test_chat.py index b21aae75..0bed93fd 100644 --- a/backend/tests/test_chat.py +++ b/backend/tests/test_chat.py @@ -30,6 +30,7 @@ from druks.harnesses.opencode import OpenCodeHarness from druks.mcp.enums import Toolkit from druks.mcp.inbound import get_druks_account_token, get_druks_mcp_server +from druks.mcp.models import McpServer from druks.models import Base from druks.redis import get_client from druks.sandbox.exceptions import IdentityDenied @@ -90,6 +91,7 @@ async def sandbox(druks_db, conversation, monkeypatch): identity={}, effort="", fast_mode=False, + timeout=600, ) monkeypatch.setattr(service, "get_agent", AsyncMock(return_value=(config, "", Toolkit.ALL))) monkeypatch.setattr(service, "sandbox_client", SimpleNamespace(set_expiry=AsyncMock())) @@ -200,17 +202,29 @@ async def attach(*, host_id): ) config = SimpleNamespace(secret_refs=[], secrets={}) first_host, _identity = await service.get_sandbox( - druks_db, conversation.account_id, config, Toolkit.ALL + druks_db, + conversation.account_id, + config=config, + allowed_tools=Toolkit.ALL, + mcp_secret_refs=[], ) second_host, _identity = await service.get_sandbox( - druks_db, second.account_id, config, Toolkit.ALL + druks_db, + second.account_id, + config=config, + allowed_tools=Toolkit.ALL, + mcp_secret_refs=[], ) login = await get_druks_account_token(druks_db, conversation.account_id, (), name="login") moved = SimpleNamespace( secret_refs=[SecretRef(name="claude_token", secret_id=login.id)], secrets={} ) moved_host, _identity = await service.get_sandbox( - druks_db, second.account_id, moved, Toolkit.ALL + druks_db, + second.account_id, + config=moved, + allowed_tools=Toolkit.ALL, + mcp_secret_refs=[], ) assert first_host.id == second_host.id @@ -461,6 +475,48 @@ async def download(**values): assert previous.deleted_at +LINEAR = { + "name": "linear", + "url": "https://mcp.linear.app/mcp", + "headers": {"Authorization": "${MCP_LINEAR_HEADER_0}"}, +} + + +@pytest.mark.parametrize( + "kind, expected", [(AccountKind.OPERATOR, [LINEAR]), (AccountKind.BOT, [])] +) +async def test_only_the_operator_reaches_the_connected_mcp_servers( + druks_db, conversation, sandbox, monkeypatch, kind, expected +): + await McpServer.create( + druks_db, + name="linear", + url="https://mcp.linear.app/mcp", + secret_headers={"Authorization": "Bearer lin_secret"}, + ) + await McpServer.create(druks_db, name="sentry", url="https://mcp.sentry.dev/mcp", is_oauth=True) + conversation.account.kind = kind + starts = [] + + async def request(self, method, **values): + if method == "start": + starts.append(values) + return {"status": "idle", "sessionId": "one"} + + async def follow_turn(session, conversation, message, bridge): + message.state = MessageState.REPLIED + await session.commit() + + monkeypatch.setattr(Bridge, "request", request) + monkeypatch.setattr(service, "follow_turn", follow_turn) + + await service.deliver_pending(druks_db, conversation) + + [refs] = [call.kwargs["mcp_secret_refs"] for call in service.get_sandbox.await_args_list] + assert len(refs) == len(expected) + assert starts[0]["mcpServers"] == expected + + @pytest.mark.parametrize( "harness, account_type, expected", [ @@ -883,7 +939,7 @@ async def test_stop_during_startup_keeps_the_sandbox_and_sends_only_the_next_mes sandbox_requests = 0 prompts = [] - async def get_sandbox(session, account_id, config, allowed_tools): + async def get_sandbox(session, account_id, *, config, allowed_tools, mcp_secret_refs): nonlocal sandbox_requests sandbox_requests += 1 await session.commit() diff --git a/backend/tests/test_whatsapp.py b/backend/tests/test_whatsapp.py index ede6148a..c56256c1 100644 --- a/backend/tests/test_whatsapp.py +++ b/backend/tests/test_whatsapp.py @@ -781,7 +781,14 @@ async def test_operator_turns_use_the_apps_prompt_and_settings_without_a_timeout resolved_config, prompt, tools = await service.get_agent(druks_db, conversation) await service.send_turn( - druks_db, conversation, message, Bridge(host), SimpleNamespace(), resolved_config, prompt + druks_db, + conversation, + message, + bridge=Bridge(host), + identity=SimpleNamespace(), + config=resolved_config, + prompt=prompt, + mcp_servers=(), ) bot = helpdesk.bot if app == "helpdesk" else Chat.bot diff --git a/docs/chat.md b/docs/chat.md index d7c2d386..238444ca 100644 --- a/docs/chat.md +++ b/docs/chat.md @@ -102,6 +102,12 @@ An operator's Chat key permits every tool of the Druks toolkit: the routes tagge tool list lists and calls only those tools. The agents of a WhatsApp number have such keys: see [WhatsApp](#whatsapp). +An operator's agent also reaches every MCP server that is enabled in +**Settings → MCP servers**, with the same credentials as a workflow agent. Chat +leaves out a server that cannot authenticate for you, such as an OAuth server +that you have not connected. A Bot's agent reaches only the Druks server. When +you enable, disable, or connect a server, the next turn replaces the sandbox. + The agent can change Druks through those tools. Chat has no permission dialog or proposal mode. Claude runs in bypass mode and cannot call `AskUserQuestion`. Codex runs in full-access mode for an operator. For a Bot it runs in read-only