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
6 changes: 3 additions & 3 deletions backend/druks/agents.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 = [
Expand All @@ -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)

Expand Down
10 changes: 9 additions & 1 deletion backend/druks/chat/bridge.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
38 changes: 36 additions & 2 deletions backend/druks/chat/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -103,15 +104,18 @@ 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."""
server = get_druks_mcp_server(allowed_tools=())
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,
Expand Down Expand Up @@ -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)
Expand All @@ -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."""
Expand Down Expand Up @@ -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
),
Expand Down
24 changes: 10 additions & 14 deletions backend/druks/harnesses/claude.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]:
Expand Down
6 changes: 1 addition & 5 deletions backend/druks/harnesses/pi.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
9 changes: 9 additions & 0 deletions backend/druks/sandbox/datastructures.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
18 changes: 14 additions & 4 deletions backend/druks/workspaces.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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
Expand Down
13 changes: 8 additions & 5 deletions backend/tests/chat/bridge.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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'}]}; },",
Expand Down Expand Up @@ -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);
Expand Down
64 changes: 60 additions & 4 deletions backend/tests/test_chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()))
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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",
[
Expand Down Expand Up @@ -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()
Expand Down
9 changes: 8 additions & 1 deletion backend/tests/test_whatsapp.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions docs/chat.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading