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
48 changes: 48 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -126,3 +126,51 @@ You stay in main thread, checking in on agents and spawning new ones as needed.
## License

This project is licensed under the [Sustainable Use License v1.0](LICENSE.md) - a source-available license that allows free use for internal business, non-commercial, and personal purposes.

## Agent tools and network isolation

Native agents use the `mainloop` MCP server, a dedicated stateless Streamable HTTP listener
on port 8002. Service `mainloop-mcp` serves port 80 at
`http://mainloop-mcp.mainloop.svc.cluster.local/mcp`; it exposes no REST API.
The Substrate egress gateway replaces the agent's literal `Authorization: Bearer
kagent-credential-injected` placeholder with its binding credential. Mainloop stores a token
hash on the binding and publishes the gateway credential under that binding id in
`kagent/mainloop-agent-tokens`. Terminal or archived bindings lose tool access.

A **NetworkPolicy-enforcing CNI is required**. The REST API has no application authorization;
its port 8000 must admit only the frontend and the Tailscale gateway. MCP port 8002 admits
only `ate-system` egress pods labelled `app: atenet-egress`. Verify those gateway pod labels
and the Tailscale gateway selector against the installation, and prove blocked connections
fail before deployment. A default Kind CNI does not provide this enforcement. The cleartext
gateway-to-Mainloop hop depends on this isolation and the pinned stock agentgateway path;
never expose the MCP Service outside the cluster.

The base includes the dedicated MCP container, Service and ingress policy. Cross-namespace
bootstrap resources are separately rendered with `k8s/integrations/kagent`: the empty token
Secret, name-scoped Role/RoleBinding, and RemoteMCPServer. Configure GitOps to preserve the
Secret's runtime-managed data. Each native AgentTemplate must bind that RemoteMCPServer.
Configure `KAGENT_MAIN_AGENT` (default `mainloop-main`) as a dedicated Claude Agent whose
Harness has `sessionIdleTTL: 0s`; child agents retain their own TTLs. Agent templates and
harnesses remain owned by the kagent installation. Gateway port 8083 must admit only Mainloop,
and TaskStore must admit only actors, using installation-specific policies.

The supported gateway is stock Substrate v0.3.0-alpha3 using agentgateway revision
`50999825cb55904801f7fd6b18b865179d0d50c4`, image digest
`sha256:f1907a50b2e74a071da53fcd2008d585b6a63d31b4e1ba3ee46cf22b342cf04b`.
That dataplane injects complete header values on HTTP and HTTPS when a configured placeholder
header is present. The credential provider must authorize the actor's atespace to read the
`kagent` namespace containing `mainloop-agent-tokens`; Mainloop's name-scoped writer Role does
not grant the provider that access. Keep this version pinned and repeat the gateway regression
proof on every Substrate/agentgateway upgrade: HTTP injection is not a portable guarantee of
other dataplanes. The proposed atenet cleartext allowlist patch is parked and is not required.

The companion's hash-only live gateway proof established HTTP/HTTPS injection and placeholder
isolation on that pinned stock dataplane. Joint Mainloop/native-agent MCP proof and blocked
connection/CNI verification remain pending. The Secret holds the complete `Bearer ml_…` header
value; NetworkPolicy remains required even though the gateway proof passed.

Slice a2 requires a **fresh database**. There is no supported upgrade from earlier slices or
credential-less kagent bindings. Reset dev/spike data before using this channel. Main and child
Sessions start with credential references; the main thread uses its dedicated main Agent.
Earlier native session history is not carried into a2. Mainloop does not migrate or reuse
pre-a2 sessions.
1 change: 1 addition & 0 deletions backend/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ dependencies = [
"dbos>=2.7.0",
"pydantic-ai[dbos]>=1.39.0",
"kubernetes>=34.1.0",
"mcp>=1.25.0",
]

[project.scripts]
Expand Down
2 changes: 0 additions & 2 deletions backend/src/mainloop/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
ConversationListResponse,
ConversationResponse,
)
from mainloop.runtime.agent_api import router as agent_api_router
from mainloop.runtime.credential_reauth_api import (
router as credential_reauth_api_router,
)
Expand Down Expand Up @@ -310,7 +309,6 @@ async def list_topics(user_id: str = Header(alias="X-User-ID", default=None)):
return out


app.include_router(agent_api_router)
app.include_router(workspace_api_router)
app.include_router(credential_reauth_api_router)
register_preview_proxy(app)
Expand Down
1 change: 1 addition & 0 deletions backend/src/mainloop/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ def shim_token_secret_name(self, atespace: str, actor: str) -> str:
# Session becomes not found and is replaced, losing its native context.
kagent_user_id: str = "mainloop"
kagent_namespace: str = "kagent"
kagent_main_agent: str = "mainloop-main"
kagent_claude_agent: str = "claude-subscription"
kagent_codex_agent: str = "codex-subscription-https"
kagent_request_timeout_seconds: float = 30.0
Expand Down
14 changes: 14 additions & 0 deletions backend/src/mainloop/db/postgres.py
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,8 @@ def _parse_json_field(value: Any) -> list | dict | None:
kagent_request_id TEXT, -- CreateSession request id of a replacement Session; NULL = derived
model TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
credential_cleanup_pending BOOLEAN NOT NULL DEFAULT FALSE,
child_start_failure TEXT,
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
-- Delivery ledger: one row per message; the A2A task is the receipt
Expand Down Expand Up @@ -1506,6 +1508,10 @@ async def archive_sessions(
session_ids,
parent_session_id,
)
from mainloop.runtime.agent_credentials import revoke

for row in rows:
await revoke(row["id"])
return [r["id"] for r in rows]

async def update_session(
Expand Down Expand Up @@ -1641,6 +1647,14 @@ async def update_session(
f"UPDATE sessions SET {', '.join(updates)} WHERE id = ${param_idx}",
*params,
)
if status in (
SessionStatus.COMPLETED,
SessionStatus.FAILED,
SessionStatus.CANCELLED,
):
from mainloop.runtime.agent_credentials import revoke

await revoke(session_id)

def _row_to_session(self, row: asyncpg.Record) -> Session:
result = _parse_json_field(row.get("result"))
Expand Down
175 changes: 175 additions & 0 deletions backend/src/mainloop/mcp_app.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
"""Dedicated, stateless Mainloop MCP listener; it exposes no REST routes."""

from contextlib import asynccontextmanager
from contextvars import ContextVar

from fastapi import HTTPException
from mainloop.db import db
from mainloop.runtime.agent_tools import AgentService, Ctx
from mainloop.runtime.policy import PolicyError, may_call, tools_for
from mcp.server.fastmcp import FastMCP
from mcp.server.transport_security import TransportSecuritySettings
from mcp.types import CallToolResult, TextContent, Tool
from pydantic import ValidationError
from starlette.responses import Response

from models.agent_tools import (
Delegate,
OptionalSession,
PendingDone,
Read,
Record,
Report,
RequiredSession,
ToolInput,
TopicOpen,
)

_context: ContextVar[Ctx] = ContextVar("mainloop_agent")

# Registry is shared by discovery and invocation; future approval groups can select subsets.
TOOLS = {
"whoami": (ToolInput, "Show your binding identity and role."),
"topics": (ToolInput, "List topics and pending counts."),
"topic_open": (TopicOpen, "Create or select a topic and set its status."),
"note": (Record, "Write a durable topic note."),
"decide": (Record, "Record a topic decision."),
"pending_add": (Record, "Record pending user intent."),
"pending_done": (PendingDone, "Close a pending item by id or unique prefix."),
"delegate": (Delegate, "Start a child agent for a task."),
"report": (Report, "Report your result to your parent exactly once."),
"status": (OptionalSession, "Read child state without sending a turn."),
"read": (Read, "Read stored child messages."),
"cancel": (RequiredSession, "Stop a child in your tree."),
"clear": (OptionalSession, "Archive finished children in your tree."),
}


async def invoke(
service: AgentService, ctx: Ctx, name: str, arguments: dict
) -> CallToolResult:
try:
may_call(ctx.actor, name)
body = TOOLS[name][0].model_validate(arguments)
if name == "whoami":
b = ctx.binding
result = {
"text": f"{b['role']} {b['kind']} session={b['session_id'][:8]} depth={ctx.actor.depth}",
"session_id": b["session_id"],
"role": b["role"],
"depth": ctx.actor.depth,
}
elif name in ("note", "decide", "pending_add"):
result = await service.record(
ctx,
{"note": "note", "decide": "decision", "pending_add": "pending"}[name],
body.text,
body.topic,
)
elif name == "pending_done":
result = await service.done(ctx, body.id)
else:
args = body.model_dump()
if "session" in args:
args["session_id"] = args.pop("session")
result = await getattr(service, name)(ctx, **args)
return CallToolResult(
content=[TextContent(type="text", text=result["text"])],
structuredContent=result,
)
except PolicyError as exc:
error = f"[{exc.code}] {exc.message}"
except ValidationError:
error = "[validation] invalid tool arguments"
except HTTPException as exc:
error = str(exc.detail)
if not error.startswith("["):
error = f"[{exc.status_code}] {error}"
return CallToolResult(content=[TextContent(type="text", text=error)], isError=True)


class AgentAuth:
def __init__(self, app, service):
self.app, self.service = app, service

async def __call__(self, scope, receive, send):
if scope["type"] != "http":
return await self.app(scope, receive, send)
if scope["path"] != "/mcp":
return await Response(status_code=404)(scope, receive, send)
headers = dict(scope["headers"])
scheme, _, token = (
headers.get(b"authorization", b"").decode("latin1").partition(" ")
)
try:
if scheme.lower() != "bearer" or not token:
raise HTTPException(401)
ctx = await self.service.authenticate(token)
except HTTPException:
return await Response(status_code=401)(scope, receive, send)
marker = _context.set(ctx)
try:
await self.app(scope, receive, send)
finally:
_context.reset(marker)


def create_app(service: AgentService | None = None):
managed = service is None
if service is None:
from mainloop.runtime.delegation import PgStore

service = AgentService(PgStore())
server = FastMCP(
"mainloop",
stateless_http=True,
json_response=True,
transport_security=TransportSecuritySettings(
allowed_hosts=[
"mainloop-mcp.mainloop.svc.cluster.local",
"mainloop-mcp.mainloop.svc.cluster.local:*",
"localhost:*",
"127.0.0.1:*",
"testserver",
],
allowed_origins=[],
),
)

@server._mcp_server.list_tools()
async def list_tools():
return [
Tool(
name=name,
description=description,
inputSchema=model.model_json_schema(),
)
for name, (model, description) in TOOLS.items()
if name in tools_for(_context.get().actor)
]

@server._mcp_server.call_tool(validate_input=False)
async def call_tool(name, arguments):
return await invoke(service, _context.get(), name, arguments)

http_app = server.streamable_http_app()

@asynccontextmanager
async def lifespan(app):
if managed:
await db.connect()
try:
async with server.session_manager.run():
yield
finally:
if managed:
from mainloop.runtime.native_sessions import close_client

await close_client()
await db.disconnect()

http_app.router.lifespan_context = lifespan
return AgentAuth(http_app, service)


app = create_app()
99 changes: 99 additions & 0 deletions backend/src/mainloop/runtime/agent_credentials.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
"""Publish per-binding gateway credentials. Never log Secret bodies or API exceptions."""

import asyncio
import base64
import logging

from kubernetes import client, config
from kubernetes.client.exceptions import ApiException
from mainloop.config import settings
from mainloop.runtime.agent_identity import token_for
from mainloop.runtime.kagent_client import SessionCredential

SECRET_NAME = (
"mainloop-agent-tokens" # nosec B105 - Kubernetes object name, not a credential
)
logger = logging.getLogger(__name__)
MCP_ORIGIN = "http://mainloop-mcp.mainloop.svc.cluster.local"


def credential_value(binding_id: str) -> str:
return "Bearer " + token_for(binding_id)


class CredentialStore:
def __init__(self, api=None):
self.api = api

def _api(self):
if self.api is None:
config.load_incluster_config()
self.api = client.CoreV1Api()
return self.api

def _write(self, binding_id: str, value: str | None):
# A merge patch touches one key only, so concurrent bindings cannot overwrite others.
data = base64.b64encode(value.encode()).decode() if value is not None else None
try:
self._api().patch_namespaced_secret(
SECRET_NAME,
settings.kagent_namespace,
{"data": {binding_id: data}},
_request_timeout=(5, 15),
)
except ApiException:
raise RuntimeError("could not update Mainloop agent credential") from None

async def publish(self, binding_id: str) -> SessionCredential:
await asyncio.to_thread(self._write, binding_id, credential_value(binding_id))
return SessionCredential(
origin=MCP_ORIGIN,
header="Authorization",
secret_name=SECRET_NAME,
secret_key=binding_id,
)

async def remove(self, binding_id: str):
await asyncio.to_thread(self._write, binding_id, None)


credentials = CredentialStore()


async def revoke(binding_id: str):
from mainloop.db import db

async with db.connection() as conn:
revoked = await conn.fetchval(
"UPDATE native_bindings SET token_hash=NULL, credential_cleanup_pending=TRUE WHERE session_id=$1 AND role IN ('main','child') RETURNING session_id",
binding_id,
)
if revoked:
await _cleanup(binding_id)


async def _cleanup(binding_id: str):
from mainloop.db import db

try:
await credentials.remove(binding_id)
except Exception:
# Auth is already revoked. Do not lose a report or cancellation to a Kubernetes outage.
logger.warning("agent credential cleanup pending for binding %s", binding_id)
return
async with db.connection() as conn:
await conn.execute(
"UPDATE native_bindings SET credential_cleanup_pending=FALSE WHERE session_id=$1",
binding_id,
)


async def reconcile_cleanup():
from mainloop.db import db

async with db.connection() as conn:
rows = await conn.fetch(
"SELECT session_id FROM native_bindings WHERE credential_cleanup_pending=TRUE"
)
for row in rows:
await _cleanup(row["session_id"])
Loading
Loading