fix(security): F09 (Astra P1) — CRM-/MCP-Tools delegieren mit HMAC-Token statt toter Header
Check Cross-Plugin Imports / check (push) Has been cancelled
Check Cross-Plugin Imports / check (push) Has been cancelled
Vorher: Zwei generische CRM-API-Tools (ai_assistant/crm_api_tool, mcp_server/tool_definitions) sendeten X-Internal-Call/X-Tenant-Id/ X-User-Id-Header — die geschuetzte API akzeptiert diese nicht als Authentisierung (Astra-Repro: "Not authenticated"). Der vorhandene Delegationstoken-Code (app/core/delegation_token.py, HMAC-SHA256, max 60s) war komplett unverbunden (0 Aufrufer). Im Worker zeigte der lokale Default-Host zudem auf den Worker selbst. Fix: - get_current_user akzeptiert X-Delegation-Token: HMAC-verifiziert, baut den User-Kontext mit den ECHTEN Rechten des Users auf (get_cached_permissions + RLS-Kontext) — keine Sonderrechte - CSRF-Middleware skippt Delegations-Header (browsers never attach them cross-site — gleiche Begruendung wie Bearer) - Gemeinsamer Helper _make_internal_api_request in crm_api_tool: erstellt pro Request ein 60s-Delegationstoken, sendet es als X-Delegation-Token; MCP-Tool delegiert an denselben Helper (Astra: beide Implementierungen konsolidieren) - _get_base_url: INTERNAL_API_URL-Override — Compose setzt fuer den Worker http://crm_app:8000 (127.0.0.1 zeigte im Worker auf sich selbst) Abnahme (Astra): Dieselbe Fachaktion ist fuer denselben Benutzer ueber UI und Agent gleichermaassen erlaubt oder gesperrt — die Tools laufen jetzt mit den echten User-Rechten durch denselben Auth-Pfad. Das Audit-Naming (delegated_by) folgt mit dem transparency-Update. Verifikation: test_api_tokens (inkl. 6 Delegations-Tests) + test_agent_loop + test_s1_security_guards 49/49, Syntax + ruff clean.
This commit is contained in:
@@ -83,6 +83,14 @@ class CSRFMiddleware(BaseHTTPMiddleware):
|
|||||||
if auth_header.startswith("Bearer "):
|
if auth_header.startswith("Bearer "):
|
||||||
return await call_next(request)
|
return await call_next(request)
|
||||||
|
|
||||||
|
# F09 (Astra P1): internal delegation calls carry a short-lived
|
||||||
|
# HMAC-signed X-Delegation-Token (created server-side by the
|
||||||
|
# CRM-API tool / MCP server, max 60 seconds) — CSRF-immune for the
|
||||||
|
# same reason as Bearer: browsers never attach this header to a
|
||||||
|
# cross-site request.
|
||||||
|
if request.headers.get("x-delegation-token"):
|
||||||
|
return await call_next(request)
|
||||||
|
|
||||||
if request.method in self.UNSAFE_METHODS:
|
if request.method in self.UNSAFE_METHODS:
|
||||||
# 1. Origin header check
|
# 1. Origin header check
|
||||||
origin = request.headers.get("origin")
|
origin = request.headers.get("origin")
|
||||||
|
|||||||
+55
@@ -74,7 +74,62 @@ async def get_current_user(
|
|||||||
|
|
||||||
Returns session data dict with user_id, tenant_id, email, name, role,
|
Returns session data dict with user_id, tenant_id, email, name, role,
|
||||||
and resolved permissions from Redis cache.
|
and resolved permissions from Redis cache.
|
||||||
|
|
||||||
|
F09 (Astra P1): also accepts a short-lived HMAC-signed delegation
|
||||||
|
token (X-Delegation-Token header) for INTERNAL calls made on behalf
|
||||||
|
of a user — e.g. the generic CRM-API tool used by AI agents and the
|
||||||
|
MCP server. Previously those tools sent unauthenticated
|
||||||
|
X-Internal-Call headers that the protected API never accepted.
|
||||||
"""
|
"""
|
||||||
|
# F09: internal delegation path — HMAC-signed, max 60 seconds
|
||||||
|
delegation_header = request.headers.get("X-Delegation-Token", "")
|
||||||
|
if delegation_header:
|
||||||
|
from app.core.delegation_token import verify_delegation_token
|
||||||
|
|
||||||
|
payload = verify_delegation_token(delegation_header)
|
||||||
|
if payload is None:
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_401_UNAUTHORIZED,
|
||||||
|
detail={"detail": "Invalid or expired delegation token", "code": "delegation_invalid"},
|
||||||
|
)
|
||||||
|
deleg_user_id = str(payload.get("user_id", ""))
|
||||||
|
deleg_tenant_id = str(payload.get("tenant_id", ""))
|
||||||
|
if not deleg_user_id or not deleg_tenant_id:
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_401_UNAUTHORIZED,
|
||||||
|
detail={"detail": "Delegation token missing user/tenant", "code": "delegation_invalid"},
|
||||||
|
)
|
||||||
|
tenant_id = uuid.UUID(deleg_tenant_id)
|
||||||
|
user_id = uuid.UUID(deleg_user_id)
|
||||||
|
await set_tenant_context(db, tenant_id)
|
||||||
|
from app.core.permissions import get_cached_permissions
|
||||||
|
|
||||||
|
resolved = await get_cached_permissions(db, redis, user_id, tenant_id)
|
||||||
|
from sqlalchemy import select as _select
|
||||||
|
|
||||||
|
from app.models.group import UserGroup
|
||||||
|
|
||||||
|
groups_q = await db.execute(
|
||||||
|
_select(UserGroup.group_id)
|
||||||
|
.where(UserGroup.user_id == user_id)
|
||||||
|
.where(UserGroup.tenant_id == tenant_id)
|
||||||
|
)
|
||||||
|
group_ids = [row[0] for row in groups_q]
|
||||||
|
is_admin = bool(resolved.get("is_system_admin", False))
|
||||||
|
await set_user_context(db, user_id, group_ids, is_admin)
|
||||||
|
return {
|
||||||
|
"user_id": deleg_user_id,
|
||||||
|
"tenant_id": deleg_tenant_id,
|
||||||
|
"email": "", # not needed for permission decisions
|
||||||
|
"name": "delegated",
|
||||||
|
"role": "",
|
||||||
|
"permissions": resolved.get("permissions", []),
|
||||||
|
"denied_permissions": resolved.get("denied", []),
|
||||||
|
"field_permissions": resolved.get("field_permissions", {}),
|
||||||
|
"is_system_admin": is_admin,
|
||||||
|
"delegated_by": payload.get("agent_id", ""),
|
||||||
|
}
|
||||||
|
|
||||||
settings = get_settings()
|
settings = get_settings()
|
||||||
|
|
||||||
session_id = request.cookies.get(settings.session_cookie_name)
|
session_id = request.cookies.get(settings.session_cookie_name)
|
||||||
|
|||||||
@@ -20,12 +20,56 @@ _openapi_cache: dict[str, Any] | None = None
|
|||||||
|
|
||||||
|
|
||||||
def _get_base_url() -> str:
|
def _get_base_url() -> str:
|
||||||
"""Get the internal base URL for API calls."""
|
"""Get the internal base URL for API calls.
|
||||||
|
|
||||||
|
F09 (Astra P1): inside the API container this is 127.0.0.1; the
|
||||||
|
WORKER must reach the API service instead - 127.0.0.1 there points
|
||||||
|
at the worker itself. INTERNAL_API_URL overrides (compose sets it
|
||||||
|
to http://crm_app:PORT for the worker service).
|
||||||
|
"""
|
||||||
import os
|
import os
|
||||||
|
|
||||||
|
override = os.environ.get("INTERNAL_API_URL")
|
||||||
|
if override:
|
||||||
|
return override.rstrip("/")
|
||||||
port = os.environ.get("PORT", "8000")
|
port = os.environ.get("PORT", "8000")
|
||||||
return f"http://127.0.0.1:{port}"
|
return f"http://127.0.0.1:{port}"
|
||||||
|
|
||||||
|
|
||||||
|
async def _make_internal_api_request(method, path, tenant_id, user_id, body=None):
|
||||||
|
"""F09 (Astra P1): authenticated internal API request.
|
||||||
|
|
||||||
|
Sends a short-lived HMAC-signed delegation token (max 60 s) instead
|
||||||
|
of the previous unauthenticated X-Internal-Call headers that the
|
||||||
|
protected API never accepted. The token acts ON BEHALF OF the user
|
||||||
|
- their real permissions apply (RBAC + RLS), no special rights.
|
||||||
|
"""
|
||||||
|
from app.core.delegation_token import create_delegation_token
|
||||||
|
|
||||||
|
token = create_delegation_token(
|
||||||
|
user_id=str(user_id),
|
||||||
|
tenant_id=str(tenant_id),
|
||||||
|
agent_id="crm-api-tool",
|
||||||
|
)
|
||||||
|
headers = {
|
||||||
|
"Content-Type": "application/json",
|
||||||
|
"X-Delegation-Token": token,
|
||||||
|
}
|
||||||
|
url = f"{_get_base_url()}{path}"
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
if method == "GET":
|
||||||
|
return await client.get(url, headers=headers, timeout=30.0)
|
||||||
|
if method == "POST":
|
||||||
|
return await client.post(url, headers=headers, json=body, timeout=30.0)
|
||||||
|
if method == "PATCH":
|
||||||
|
return await client.patch(url, headers=headers, json=body, timeout=30.0)
|
||||||
|
if method == "PUT":
|
||||||
|
return await client.put(url, headers=headers, json=body, timeout=30.0)
|
||||||
|
if method == "DELETE":
|
||||||
|
return await client.delete(url, headers=headers, timeout=30.0)
|
||||||
|
raise ValueError(f"Unsupported method: {method}")
|
||||||
|
|
||||||
|
|
||||||
async def get_openapi_spec() -> dict[str, Any]:
|
async def get_openapi_spec() -> dict[str, Any]:
|
||||||
"""Get the CRM OpenAPI spec, cached."""
|
"""Get the CRM OpenAPI spec, cached."""
|
||||||
global _openapi_cache
|
global _openapi_cache
|
||||||
@@ -101,34 +145,15 @@ async def call_crm_api_handler(arguments: dict[str, Any], context: dict[str, Any
|
|||||||
path = "/" + path
|
path = "/" + path
|
||||||
|
|
||||||
try:
|
try:
|
||||||
base_url = _get_base_url()
|
|
||||||
|
|
||||||
# Get user context for auth
|
# Get user context for auth
|
||||||
tenant_id = context.get("tenant_id", "")
|
tenant_id = context.get("tenant_id", "")
|
||||||
user_id = context.get("user_id", "")
|
user_id = context.get("user_id", "")
|
||||||
|
|
||||||
# Create a DB session to resolve a valid session token for this user
|
# F09 (Astra P1): authenticated request via short-lived delegation
|
||||||
# We'll use internal service-level auth bypass
|
# token - the acting user's real permissions apply.
|
||||||
headers = {
|
resp = await _make_internal_api_request(
|
||||||
"Content-Type": "application/json",
|
method, path, tenant_id=tenant_id, user_id=user_id, body=body
|
||||||
"X-Internal-Call": "true",
|
)
|
||||||
"X-Tenant-Id": str(tenant_id),
|
|
||||||
"X-User-Id": str(user_id),
|
|
||||||
}
|
|
||||||
|
|
||||||
async with httpx.AsyncClient() as client:
|
|
||||||
if method == "GET":
|
|
||||||
resp = await client.get(f"{base_url}{path}", headers=headers, timeout=30.0)
|
|
||||||
elif method == "POST":
|
|
||||||
resp = await client.post(f"{base_url}{path}", headers=headers, json=body, timeout=30.0)
|
|
||||||
elif method == "PATCH":
|
|
||||||
resp = await client.patch(f"{base_url}{path}", headers=headers, json=body, timeout=30.0)
|
|
||||||
elif method == "PUT":
|
|
||||||
resp = await client.put(f"{base_url}{path}", headers=headers, json=body, timeout=30.0)
|
|
||||||
elif method == "DELETE":
|
|
||||||
resp = await client.delete(f"{base_url}{path}", headers=headers, timeout=30.0)
|
|
||||||
else:
|
|
||||||
return json.dumps({"error": f"Unsupported method: {method}"})
|
|
||||||
|
|
||||||
# Return response body (truncated if too large)
|
# Return response body (truncated if too large)
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -12,7 +12,6 @@ import os
|
|||||||
import uuid
|
import uuid
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
import httpx
|
|
||||||
from sqlalchemy.ext.asyncio import AsyncSession
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
from app.plugins.builtins.mcp_server.schemas import McpToolDefinition, McpToolParameter
|
from app.plugins.builtins.mcp_server.schemas import McpToolDefinition, McpToolParameter
|
||||||
@@ -95,30 +94,16 @@ async def _handler_call_crm_api(db: AsyncSession, arguments: dict[str, Any], con
|
|||||||
path = "/" + path
|
path = "/" + path
|
||||||
|
|
||||||
try:
|
try:
|
||||||
base_url = _get_base_url()
|
|
||||||
tenant_id = context.get("tenant_id", "")
|
tenant_id = context.get("tenant_id", "")
|
||||||
user_id = context.get("user_id", "")
|
user_id = context.get("user_id", "")
|
||||||
|
|
||||||
headers = {
|
# F09 (Astra P1): authenticated request via short-lived delegation
|
||||||
"Content-Type": "application/json",
|
# token - the MCP session user's real permissions apply.
|
||||||
"X-Internal-Call": "true",
|
from app.plugins.builtins.ai_assistant.crm_api_tool import _make_internal_api_request
|
||||||
"X-Tenant-Id": str(tenant_id),
|
|
||||||
"X-User-Id": str(user_id),
|
|
||||||
}
|
|
||||||
|
|
||||||
async with httpx.AsyncClient() as client:
|
resp = await _make_internal_api_request(
|
||||||
if method == "GET":
|
method, path, tenant_id=str(tenant_id), user_id=str(user_id), body=body
|
||||||
resp = await client.get(f"{base_url}{path}", headers=headers, timeout=30.0)
|
)
|
||||||
elif method == "POST":
|
|
||||||
resp = await client.post(f"{base_url}{path}", headers=headers, json=body, timeout=30.0)
|
|
||||||
elif method == "PATCH":
|
|
||||||
resp = await client.patch(f"{base_url}{path}", headers=headers, json=body, timeout=30.0)
|
|
||||||
elif method == "PUT":
|
|
||||||
resp = await client.put(f"{base_url}{path}", headers=headers, json=body, timeout=30.0)
|
|
||||||
elif method == "DELETE":
|
|
||||||
resp = await client.delete(f"{base_url}{path}", headers=headers, timeout=30.0)
|
|
||||||
else:
|
|
||||||
return {"error": f"Unsupported method: {method}"}
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
resp_data = resp.json()
|
resp_data = resp.json()
|
||||||
|
|||||||
@@ -123,6 +123,7 @@ services:
|
|||||||
WORKER_DATABASE_URL: ${WORKER_DATABASE_URL:-postgresql+asyncpg://crm_worker:${DB_PASSWORD}@postgres:5432/${POSTGRES_DB:-crm_db}}
|
WORKER_DATABASE_URL: ${WORKER_DATABASE_URL:-postgresql+asyncpg://crm_worker:${DB_PASSWORD}@postgres:5432/${POSTGRES_DB:-crm_db}}
|
||||||
MIGRATION_DATABASE_URL: ${MIGRATION_DATABASE_URL:-postgresql+asyncpg://crm_user:${DB_PASSWORD}@postgres:5432/${POSTGRES_DB:-crm_db}}
|
MIGRATION_DATABASE_URL: ${MIGRATION_DATABASE_URL:-postgresql+asyncpg://crm_user:${DB_PASSWORD}@postgres:5432/${POSTGRES_DB:-crm_db}}
|
||||||
REDIS_URL: ${REDIS_URL:-redis://default:${REDIS_PASSWORD}@redis:6379/0}
|
REDIS_URL: ${REDIS_URL:-redis://default:${REDIS_PASSWORD}@redis:6379/0}
|
||||||
|
INTERNAL_API_URL: http://crm_app:8000
|
||||||
SECRET_KEY: ${SECRET_KEY:?SECRET_KEY is required}
|
SECRET_KEY: ${SECRET_KEY:?SECRET_KEY is required}
|
||||||
FRONTEND_URL: ${FRONTEND_URL}
|
FRONTEND_URL: ${FRONTEND_URL}
|
||||||
ENVIRONMENT: ${ENVIRONMENT:-production}
|
ENVIRONMENT: ${ENVIRONMENT:-production}
|
||||||
|
|||||||
Reference in New Issue
Block a user