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
22 changes: 14 additions & 8 deletions backend/druks/agents.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,17 +59,18 @@ async def _runner(
# The agent always runs in a Workspace. A warm run attaches the run's held VM; the
# rest get a fresh ephemeral VM. Either way workflow.get_workspace() turns the VM into
# the runner — fresh per call, so nothing (connection or credential) is held across steps.
if host_id:
vm = sandbox_client.attach(host_id=host_id)
elif refs and (
identity := await SandboxIdentity.lookup(
identity = None
if refs and not host_id:
identity = await SandboxIdentity.lookup(
session,
account_id=workflow.account_id,
run_id=workflow_id,
scoped_to=step,
secret_refs=refs,
)
):
if host_id:
vm = sandbox_client.attach(host_id=host_id)
elif identity:
# A crashed attempt left its box behind. Its identity finds it again.
vm = sandbox_client.resume(host_id=identity.host_id)
else:
Expand Down Expand Up @@ -146,6 +147,9 @@ class Agent:
# ``include_plugins=False`` skips the operator's plugin state for prompts
# that hit no MCP server.
include_plugins: bool = True
# ``include_mcp=False`` gives the call no MCP server and its sandbox no
# server entry, for an agent that reads untrusted content.
include_mcp: bool = True
# ``id`` is the agent's durable key (settings, timeline, registry, step name):
# ``<app>.<attribute>`` for an agent declared on an App, or the explicit ``id=``
# of a standalone agent (a test, a one-off). ``app`` is the owning App's name,
Expand Down Expand Up @@ -334,9 +338,11 @@ async def _run(
# same servers.
subject = await workflow.subject
workspace_class = workflow.workspace_class
mcp_servers, mcp_refs = await workspace_class.get_all_mcp_servers(
session, subject, workflow.account_id
)
mcp_servers, mcp_refs = (), []
if self.include_mcp:
mcp_servers, mcp_refs = await workspace_class.get_all_mcp_servers(
session, subject, workflow.account_id
)
refs = [
*config.secret_refs,
*(
Expand Down
25 changes: 21 additions & 4 deletions backend/druks/apps/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -391,6 +391,12 @@ def operations(cls) -> "dict[str, Operation]":
name = getattr(route, "operation_id", "") or ""
if not name:
continue
if name.startswith(f"{cls.name}_"):
raise AppRouteConflict(
f"app {cls.name!r} declares operation {name!r}. Druks adds the app "
f"name to every operation id. Declare "
f"{name.removeprefix(f'{cls.name}_')!r}."
)
# An APIRoute's own path already carries its router's prefix.
path = f"/api/{cls.name}{getattr(route, 'path', '')}"
if name in found:
Expand Down Expand Up @@ -443,8 +449,12 @@ def _page_endpoint(
):
"""``wraps`` keeps the page function's signature, so FastAPI still
validates every route parameter. A page whose subject is missing answers
an empty state that links ``back``."""
from druks.ui import EmptyState, Link, Page
an empty state that links ``back``. The page serializes with where the
work on each of its subjects stands, from one read."""
from fastapi.responses import JSONResponse

from druks.durable import reads
from druks.ui import Action, EmptyState, GateControls, Link, Page, SubjectStatus

controls = []
if back and back is not declaration:
Expand All @@ -467,11 +477,18 @@ async def read_page(**parameters):
f"it answered with {type(page).__name__}, not a Page",
)
try:
for action in page.iter_actions():
for action in page.iter_parts(Action):
action.check_operation(cls.name, operations)
except ValueError as error:
raise PageContractError(cls.name, declaration.name, str(error)) from error
return page
subjects = [
(part.subject.subject_type, part.subject.subject_id)
for part in page.iter_parts(SubjectStatus, GateControls)
]
statuses = await reads.get_statuses_for_subjects(db_session(), subjects)
return JSONResponse(
page.model_dump(mode="json", by_alias=True, context={"statuses": statuses})
)

return read_page

Expand Down
17 changes: 17 additions & 0 deletions backend/druks/cli.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import argparse

from .database import make_app_migration, run_migrations
from .exceptions import DruksError
from .settings import ensure_data_dirs, load_settings, setup_logging


Expand All @@ -14,6 +15,11 @@ def main() -> None:
)
makemigrations.add_argument("app", help="The installed app's name.")
makemigrations.add_argument("-m", "--message", default="", help="Revision message (slug).")
check_app = subparsers.add_parser(
"check-app",
help="Load one installed app and check its contracts. Needs no database.",
)
check_app.add_argument("app", help="The installed app's name.")
doctor_parser = subparsers.add_parser(
"doctor",
help=(
Expand Down Expand Up @@ -118,6 +124,17 @@ def main() -> None:
print(f"Next: cd {target.name} && uv sync && uv run pytest")
return

# An app's CI checks it with no configured install.
if args.command == "check-app":
from .apps.loader import load_app

try:
load_app(args.app).routers()
except DruksError as error:
raise SystemExit(f"druks check-app: {error}") from error
print(f"{args.app}: ok")
return

settings = load_settings()
setup_logging(settings)
ensure_data_dirs(settings)
Expand Down
2 changes: 1 addition & 1 deletion backend/druks/contrib/software_factory/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -274,7 +274,7 @@ async def list_work_items_history(
@work_items_router.post(
"/{ticket}/start",
status_code=status.HTTP_202_ACCEPTED,
operation_id="software_factory_start",
operation_id="start",
tags=["agent"],
responses=agent_error_responses(TicketNotFound("ENG-9999", "Linear"), TrackerNotConfigured()),
)
Expand Down
37 changes: 24 additions & 13 deletions backend/druks/durable/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,17 @@
from typing import TYPE_CHECKING, Any, Literal

from dbos import DBOS
from sqlalchemy import CheckConstraint, ForeignKey, Index, Select, String, func, select, update
from sqlalchemy import (
CheckConstraint,
ForeignKey,
Index,
Select,
String,
func,
select,
tuple_,
update,
)
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession
Expand Down Expand Up @@ -207,38 +217,39 @@ async def get_latest_for_subject(

@classmethod
async def get_latest_for_subjects(
cls, session: AsyncSession, subject_type: str, subject_ids: list[str]
) -> dict[str, "Run"]:
"""The driving run of each subject, keyed by subject id — get_latest_for_subject
for a whole board in one statement. Agent calls come with it: the status read
needs the latest agent of every running row."""
cls, session: AsyncSession, identities: list[tuple[str, str]]
) -> dict[tuple[str, str], "Run"]:
"""The driving run of each subject, keyed by its ``(subject_type, subject_id)``:
get_latest_for_subject for subjects of any type in one statement. Agent calls
come with it: the status read needs the latest agent of every running row."""
subject_type = (
workflow_status.c.attributes["subject_type"].as_string().label("subject_type")
)
subject_id = workflow_status.c.attributes["subject_id"].as_string().label("subject_id")
driving = (
select(
subject_type,
subject_id,
cls.id.label("run_id"),
func.row_number()
.over(
partition_by=subject_id,
partition_by=(subject_type, subject_id),
order_by=(cls.created_at.desc(), cls.id.desc()),
)
.label("rank"),
)
.join_from(cls, workflow_status, workflow_status.c.workflow_uuid == cls.id)
.where(
workflow_status.c.attributes["subject_type"].as_string() == subject_type,
subject_id.in_(subject_ids),
)
.where(tuple_(subject_type, subject_id).in_(identities))
.subquery()
)
stmt = (
select(driving.c.subject_id, cls)
select(driving.c.subject_type, driving.c.subject_id, cls)
.join_from(cls, driving, driving.c.run_id == cls.id)
.where(driving.c.rank == 1)
.options(selectinload(cls.agent_calls))
)
rows = await session.execute(stmt)
return {found_id: run for found_id, run in rows}
return {(found_type, found_id): run for found_type, found_id, run in rows}

@classmethod
def open_subject_ids(cls, subject_type: str) -> Select:
Expand Down
15 changes: 13 additions & 2 deletions backend/druks/durable/reads.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,8 +93,19 @@ async def get_subject_statuses(
) -> dict[str, SubjectStatus]:
"""The status of every subject on a board, keyed by subject id — one driving-run
read for the whole page."""
driving_runs = await Run.get_latest_for_subjects(session, subject_type, subject_ids)
return {subject_id: await _status(driving_runs.get(subject_id)) for subject_id in subject_ids}
statuses = await get_statuses_for_subjects(
session, [(subject_type, subject_id) for subject_id in subject_ids]
)
return {subject_id: statuses[(subject_type, subject_id)] for subject_id in subject_ids}


async def get_statuses_for_subjects(
session: AsyncSession, identities: list[tuple[str, str]]
) -> dict[tuple[str, str], SubjectStatus]:
"""The status of each ``(subject_type, subject_id)``, keyed by it — one
driving-run read for subjects of any type."""
driving_runs = await Run.get_latest_for_subjects(session, identities)
return {identity: await _status(driving_runs.get(identity)) for identity in identities}


async def get_subject_phase(
Expand Down
38 changes: 19 additions & 19 deletions backend/druks/mcp/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ def _validate_agent_tools(api: FastAPI) -> None:
# is the one view with every route's merged tags. Validation owns only the
# two demands the author owns — an explicit operation_id and a non-empty
# docstring; the app prefix is the framework's to derive, not the
# author's to repeat (see _namespace_agent_operations).
# author's to repeat (see _namespace_app_operations).
mounted_tags: set[str] = set()
bot_operations: defaultdict[str, set[str]] = defaultdict(set)
for route in iter_route_contexts(api.routes):
Expand Down Expand Up @@ -117,25 +117,25 @@ def _validate_agent_tools(api: FastAPI) -> None:


def get_tool_name(operation_id: str, tags: list[str], app_names: set[str]) -> str:
# An app-owned agent operation's tool is f"{app}_{operation_id}", so the
# author never repeats the prefix. The loader tags every app route with its
# app's name, so among an agent operation's tags the one naming an
# installed app is the owner; platform agent operations carry no such tag
# and keep their declared ids. An already-prefixed id passes through, so
# stable names like software_factory_start never double.
# An app-owned operation's id is f"{app}_{operation_id}", so the author
# never repeats the prefix. The loader tags every app route with its app's
# name, so among an operation's tags the one naming an installed app is the
# owner; platform operations carry no such tag and keep their declared ids.
# Boot refuses an app id that already carries the prefix, so a prefixed id
# here is one this derived before, and it passes through unchanged.
app = next((tag for tag in tags if tag in app_names), None)
if app and not operation_id.startswith(f"{app}_"):
return f"{app}_{operation_id}"
return operation_id


def _namespace_agent_operations(spec: dict, app_names: set[str]) -> None:
# Rename each agent operation to its tool name: the provider reads the tool
# name off the spec. The namespace is what makes the merged document's
# operation ids globally unique. A derived id that would collide with another
# route's explicit id is rejected: before this derivation the clash was
# visible in the author's code, so the framework must surface it now that it
# owns the naming.
def _namespace_app_operations(spec: dict, app_names: set[str]) -> None:
# Rename each app operation to its namespaced id: the provider reads an
# agent operation's tool name off the spec. The namespace is what makes the
# merged document's operation ids globally unique. A derived id that would
# collide with another route's explicit id is rejected: before this
# derivation the clash was visible in the author's code, so the framework
# must surface it now that it owns the naming.
existing_ids = {
op.get("operationId")
for ops in spec.get("paths", {}).values()
Expand All @@ -144,12 +144,12 @@ def _namespace_agent_operations(spec: dict, app_names: set[str]) -> None:
}
for path, operations in spec.get("paths", {}).items():
for operation in operations.values():
if not isinstance(operation, dict) or not _TOOL_TAGS & set(operation.get("tags", [])):
if not isinstance(operation, dict):
continue
operation_id = operation.get("operationId")
if not operation_id:
continue
derived = get_tool_name(operation_id, operation["tags"], app_names)
derived = get_tool_name(operation_id, operation.get("tags", []), app_names)
if derived != operation_id:
if derived in existing_ids:
raise InvalidAgentToolError(
Expand All @@ -160,7 +160,7 @@ def _namespace_agent_operations(spec: dict, app_names: set[str]) -> None:
operation["operationId"] = derived


def _install_agent_namespacing(api: FastAPI) -> None:
def _install_app_namespacing(api: FastAPI) -> None:
# The tool name comes from the spec's operation id, so the namespace must
# land on the document api.openapi() builds — not on FastAPI's cached, merged
# route contexts, which later generation silently discards. Wrap the app's
Expand All @@ -182,7 +182,7 @@ def _install_agent_namespacing(api: FastAPI) -> None:

def namespaced() -> dict:
spec = generate()
_namespace_agent_operations(spec, app_names)
_namespace_app_operations(spec, app_names)
return spec

api.openapi = namespaced
Expand All @@ -209,7 +209,7 @@ def _is_visible(context: AuthContext) -> bool:

def create_mcp_app(api: FastAPI) -> StarletteWithLifespan:
_validate_agent_tools(api)
_install_agent_namespacing(api)
_install_app_namespacing(api)
# Built directly rather than via from_fastapi, which owns the transport:
# raise_app_exceptions=False makes an app crash reach the tool as the
# app's sanitized 500, so no masking is needed and the taxonomy travels.
Expand Down
1 change: 1 addition & 0 deletions backend/druks/scaffolding/app_template/AGENTS.md-tpl
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ webhooks, and the dashboard.
## Verify

```bash
uv run druks check-app {{ name }}
uv run pytest
```

Expand Down
3 changes: 2 additions & 1 deletion backend/druks/services/__init__.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from .base import Service
from .base import Connection, Service
from .exceptions import (
OauthExchangeError,
OauthRefreshError,
Expand All @@ -8,6 +8,7 @@
from .oauth import OauthClient

__all__ = [
"Connection",
"OauthClient",
"OauthExchangeError",
"OauthRefreshError",
Expand Down
6 changes: 6 additions & 0 deletions backend/druks/services/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,12 @@ def __set_name__(self, owner: type, name: str) -> None:
def label(self) -> str:
return f"{self.owner.name}.{self.name}"

@property
def connect_url(self) -> str:
"""Where an operator connects an account to this service. The provider
sends them back to the app."""
return f"/api/oauth/{self.service.slug}/connect?next=/{self.owner.name}"

async def list_for_account(self, account_id: str) -> list[Connection]:
return [
Connection(self.service, row)
Expand Down
2 changes: 2 additions & 0 deletions backend/druks/ui/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
Section,
Stack,
StatusValue,
SubjectStatus,
Table,
TableColumn,
TableRow,
Expand Down Expand Up @@ -101,6 +102,7 @@
"SelectField",
"Stack",
"StatusValue",
"SubjectStatus",
"Table",
"TableColumn",
"TableRow",
Expand Down
Loading
Loading