Files
leocrm/app/main.py
T
Agent Zero 727d86614e Security fixes: P0-P2 complete (22 fixes)
P0 (7): Auth-bypass removed, migrations fixed, plugin-upload disabled, RLS FORCE+WITH CHECK, plugin double-registration fixed, persistent volume, domain removed
P1 (11): User/tenant model, Redis centralized, worker separated, transactional outbox, XSS fixed, DMS chunked streaming, permissions unified, password reset, metrics secured, config/docs fixed, cross-tenant FK
P2 (4): Contact model normalized, cross-imports reduced 94%, commands+state machines for contacts/dms/mail/calendar, SPA path-traversal

8 new migrations, 99 unit tests, 13 commands, 8 contracts, 72 files changed
2026-07-25 21:03:46 +02:00

370 lines
16 KiB
Python

"""FastAPI application - LeoCRM backend."""
from __future__ import annotations
import time
import traceback
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse
from fastapi.staticfiles import StaticFiles
from starlette.middleware.base import BaseHTTPMiddleware
import importlib
import logging
import os
logger = logging.getLogger(__name__)
from app.config import get_settings
from app.core.db import close_engine, get_engine
from app.core.middleware import CSRFMiddleware
from app.core.monitoring import record_error, record_request
from app.core.service_container import get_container
from app.plugins.registry import get_registry
from app.routes import (
addresses,
bank_accounts,
ai_copilot,
audit,
auth,
contact_folders,
contacts,
dashboard,
entity_history,
groups,
health,
import_export,
metrics,
notifications,
plugins,
roles,
tenants,
users,
user_preferences,
workflows,
currencies,
taxes,
sequences,
system_settings,
attachments,
custom_fields,
saved_filters,
)
class RequestLoggingMiddleware(BaseHTTPMiddleware):
"""Structured logging + Prometheus metrics for every HTTP request."""
async def dispatch(self, request: Request, call_next):
start_time = time.perf_counter()
method = request.method
path = request.url.path
# Extract tenant_id from session cookie if available (best-effort)
tenant_id = None
try:
response = await call_next(request)
except Exception as exc:
duration_ms = (time.perf_counter() - start_time) * 1000
tb_str = traceback.format_exc()
record_error(
event="unhandled_exception",
method=method,
path=path,
status_code=500,
error=str(exc),
traceback_str=tb_str,
tenant_id=tenant_id,
)
raise
duration_ms = (time.perf_counter() - start_time) * 1000
status_code = response.status_code
# Try to get tenant_id from response headers or request state
# (set by auth middleware/dependency — best-effort, never log credentials)
record_request(
method=method,
path=path,
status_code=status_code,
duration_ms=duration_ms,
tenant_id=tenant_id,
)
return response
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Application lifespan: startup and shutdown."""
# Initialize global Redis client (singleton)
from app.core.auth import init_redis, close_redis
from app.core.jobs import init_job_pool, close_job_pool
await init_redis()
await init_job_pool()
# Initialize service container
container = get_container()
await container.initialize()
# Initialize plugin registry and discover built-in plugins
registry = get_registry()
registry.initialize(get_engine(), app)
registry.discover_builtins()
# Install discovered builtin plugins and activate only those marked active in DB
from sqlalchemy import select as sa_select
from sqlalchemy.ext.asyncio import async_sessionmaker
from app.models.plugin import Plugin as PluginModel
from app.core.event_bus import get_event_bus
event_bus = get_event_bus()
async_session = async_sessionmaker(get_engine(), expire_on_commit=False)
async with async_session() as db:
for name in registry.resolve_load_order():
plugin = registry.get_plugin(name)
if plugin is None:
continue
# Check if plugin already in DB
result = await db.execute(
sa_select(PluginModel).where(PluginModel.name == name)
)
plugin_record = result.scalar_one_or_none()
if plugin_record is None:
# Create DB record for this builtin plugin — inactive by default (except core)
plugin_record = PluginModel(
name=name,
display_name=plugin.manifest.display_name,
version=plugin.manifest.version,
status="installed",
active=plugin.manifest.is_core, # Only core plugins auto-activate
is_core=plugin.manifest.is_core,
)
db.add(plugin_record)
await db.flush()
logger.info(f"Created plugin record: {name} (core={plugin.manifest.is_core})")
# Run migrations if not yet applied
if plugin.manifest.migrations:
try:
await registry.migration_runner.run_all_migrations(
db, name, plugin.manifest.migrations
)
except Exception as exc:
logger.error(f"Migration FAILED for {name}: {exc}")
if plugin_record.active:
logger.error(f"Deactivating plugin {name} due to migration failure")
plugin_record.active = False
plugin_record.status = "migration_failed"
continue # Skip activation if migration fails
# Only activate plugins that are marked active in DB
if not plugin_record.active:
logger.info(f"Plugin {name} is inactive — skipping activation")
continue
# Activate plugin and register routes
try:
await plugin.on_activate(db, container, event_bus)
for route_def in plugin.manifest.routes:
router_module = importlib.import_module(route_def.module)
router = getattr(router_module, route_def.router_attr)
app.include_router(router)
plugin_record.status = "active"
print(f"[STARTUP] Activated plugin: {name} ({len(plugin.manifest.routes)} routes)", flush=True)
logger.info(f"Activated plugin: {name} ({len(plugin.manifest.routes)} routes)")
except Exception as exc:
print(f"[STARTUP] Failed to activate plugin {name}: {exc}", flush=True)
logger.error(f"Failed to activate plugin {name}: {exc}")
plugin_record.active = False
plugin_record.status = "activation_failed"
await db.commit()
# Initialize permission registry with active plugin names
from app.core.permission_registry import init_permission_registry, register_plugin_permissions
active_plugin_names: set[str] = set()
async with async_session() as db:
result = await db.execute(
sa_select(PluginModel).where(PluginModel.active == True) # noqa: E712
)
for record in result.scalars().all():
active_plugin_names.add(record.name)
plugin = registry.get_plugin(record.name)
if plugin and plugin.manifest.permissions:
register_plugin_permissions(record.name, plugin.manifest.permissions)
init_permission_registry(active_plugin_names)
logger.info("Permission registry initialized with %d active plugins", len(active_plugin_names))
# Register field definitions from active plugins only
from app.core.permission_registry import get_permission_registry
for name in active_plugin_names:
plugin = registry.get_plugin(name)
if plugin:
field_defs = plugin.get_field_definitions()
if field_defs:
get_permission_registry().register_field_definitions(name, field_defs)
logger.info("Field definitions registered for %d active plugins", len(active_plugin_names))
# Seed default data (EUR currency, 19%/7% tax rates) for all tenants
from app.core.seeds import seed_default_data
async with async_session() as db:
try:
await seed_default_data(db)
await db.commit()
logger.info("Default data seeding completed")
except Exception as exc:
logger.warning(f"Default data seeding failed: {exc}")
await db.rollback()
yield
# Shutdown: close global Redis and ARQ pool
await close_job_pool()
await close_redis()
await close_engine()
def create_app() -> FastAPI:
"""Create and configure the FastAPI application."""
settings = get_settings()
app = FastAPI(
title="LeoCRM",
version="1.0.0",
lifespan=lifespan,
description=(
"LeoCRM — Self-hosted CRM system for small sales teams.\n\n"
"## Authentication\n"
"All endpoints (except `/api/v1/auth/login` and `/api/v1/health`) require "
"a valid session cookie. Obtain a session by calling `POST /api/v1/auth/login`.\n\n"
"## Multi-Tenancy\n"
"All data is tenant-scoped. The tenant context is set automatically from the "
"authenticated session.\n\n"
"## Plugins\n"
"LeoCRM uses a plugin architecture. Core routes are always available. "
"Plugin routes (DMS, Mail, Calendar, Automation, etc.) are registered when "
"the corresponding plugin is active."
),
openapi_tags=[
{"name": "health", "description": "Health check endpoints — no auth required."},
{"name": "metrics", "description": "Prometheus metrics endpoint."},
{"name": "auth", "description": "Authentication: login, logout, session, password reset."},
{"name": "users", "description": "User management: CRUD, user-tenant assignments."},
{"name": "roles", "description": "Role management and permission definitions."},
{"name": "groups", "description": "User groups for contact assignment and filtering."},
{"name": "tenants", "description": "Tenant management and tenant switching."},
{"name": "notifications", "description": "User notifications and notification preferences."},
{"name": "contacts", "description": "Contact CRUD, contact persons, FTS search, export, soft-delete."},
{"name": "contact-folders", "description": "Contact folder management for organization."},
{"name": "entity-history", "description": "Audit trail and entity change history."},
{"name": "import-export", "description": "Bulk import and export of contacts and data."},
{"name": "plugins", "description": "Plugin management: list, install, activate, deactivate."},
{"name": "ai-copilot", "description": "AI copilot: chat, suggestions, conversation history."},
{"name": "workflows", "description": "Workflow definitions, instances, and execution."},
{"name": "user-preferences", "description": "Per-user preference settings."},
{"name": "currencies", "description": "Currency management for multi-currency support."},
{"name": "taxes", "description": "Tax rate management (VAT, sales tax)."},
{"name": "sequences", "description": "Number sequence management for invoices, quotes, etc."},
{"name": "system-settings", "description": "Tenant-level system settings (company info, theme)."},
{"name": "attachments", "description": "File attachments for contacts, companies, and entities."},
{"name": "addresses", "description": "Address management for contacts and companies."},
{"name": "bank_accounts", "description": "Bank account management for tenant."},
{"name": "audit", "description": "Audit log queries and compliance reporting."},
{"name": "automation", "description": "Automation engine: agents, automations, cron jobs, execution logs."},
{"name": "agents", "description": "AI agent definitions and agent runner endpoints."},
{"name": "dms", "description": "Document Management System: files, folders, sources, sharing."},
{"name": "mail", "description": "Email integration: IMAP accounts, folders, messages, send."},
{"name": "calendar", "description": "Calendar management: appointments, resources, recurrence."},
{"name": "search", "description": "Unified search across contacts, documents, emails, etc."},
{"name": "reports", "description": "Report generator: templates, rendering, scheduled reports."},
{"name": "entity-links", "description": "Entity linking: connect contacts to DMS documents and other entities."},
{"name": "kommunikation", "description": "Unified messaging: conversations, participants, messages."},
{"name": "ai-proactive", "description": "Proactive AI: insights, alerts, and recommendations."},
{"name": "ai-assistant", "description": "AI assistant: chat completions, tool calling, context awareness."},
{"name": "tags", "description": "Tag management: create, assign, search tags across entities."},
{"name": "permissions", "description": "Permission management: roles, field-level permissions, sharing."},
{"name": "public-share", "description": "Public sharing endpoints — no auth required, token-based access."},
],
)
app.add_middleware(
CORSMiddleware,
allow_origins=settings.cors_origin_list,
allow_credentials=True,
allow_methods=["GET", "POST", "PATCH", "PUT", "DELETE", "OPTIONS"],
allow_headers=["Authorization", "Content-Type", "X-CSRF-Token", "X-Tenant-ID"],
max_age=3600,
)
app.add_middleware(CSRFMiddleware)
app.add_middleware(RequestLoggingMiddleware)
app.include_router(health.router)
app.include_router(metrics.router)
app.include_router(auth.router)
app.include_router(users.router)
app.include_router(roles.router)
app.include_router(groups.router)
app.include_router(tenants.router)
app.include_router(notifications.router)
app.include_router(contacts.router)
app.include_router(contact_folders.router)
app.include_router(dashboard.router)
app.include_router(entity_history.router)
app.include_router(import_export.router)
app.include_router(plugins.router)
app.include_router(ai_copilot.router)
app.include_router(workflows.router)
app.include_router(user_preferences.router)
app.include_router(currencies.router)
app.include_router(taxes.router)
app.include_router(sequences.router)
app.include_router(system_settings.router)
app.include_router(attachments.router)
app.include_router(addresses.router)
app.include_router(bank_accounts.router)
app.include_router(audit.router)
app.include_router(custom_fields.router)
app.include_router(saved_filters.router)
# ── Plugin routes are registered in lifespan() after activation status is loaded ──
# Do NOT register plugin routes here — lifespan() handles it for active plugins only
# ── Serve frontend static files (SPA) ──────────────────────────────
# Mount built frontend assets (JS, CSS, images)
frontend_dist = os.path.join(os.path.dirname(__file__), "..", "frontend", "dist")
if os.path.isdir(frontend_dist):
# Serve static assets at /assets
assets_path = os.path.join(frontend_dist, "assets")
if os.path.isdir(assets_path):
app.mount("/assets", StaticFiles(directory=assets_path), name="assets")
# SPA catch-all: serve index.html for all non-API routes
@app.get("/{full_path:path}")
async def spa_spa(full_path: str):
"""Serve index.html for all non-API routes (SPA fallback)."""
# Don't intercept API routes
if full_path.startswith(("api/", "docs", "openapi", "redoc")):
raise HTTPException(status_code=404, detail="Not Found")
# Block path traversal and system file access
blocked_prefixes = ("var/log/", "error/", "error_log", "var/", "etc/", "proc/", "sys/")
if full_path.startswith(blocked_prefixes) or ".." in full_path:
raise HTTPException(status_code=404, detail="Not Found")
index_path = os.path.join(frontend_dist, "index.html")
if os.path.isfile(index_path):
return FileResponse(index_path)
raise HTTPException(status_code=404, detail="Frontend not built")
return app
app = create_app()