Files
leocrm/app/routes/plugins.py
T

544 lines
18 KiB
Python

"""Plugin routes — list, install, activate, deactivate, uninstall, manifest schema, config, upload, install-url."""
from __future__ import annotations
import importlib
import logging
import os
import re
import shutil
import tempfile
import zipfile
from pathlib import Path
from typing import Any
import httpx
from fastapi import APIRouter, Depends, File, HTTPException, Query, UploadFile
from pydantic import BaseModel
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.db import get_db
from app.deps import require_permission, require_admin
from app.plugins.base import BasePlugin
from app.plugins.manifest import PluginManifest
from app.plugins.migration_runner import MigrationValidationError
from app.services.plugin_service import get_plugin_service
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/v1/plugins", tags=["plugins"])
class PluginConfigUpdate(BaseModel):
"""Request body for updating plugin configuration."""
config: dict
class PluginUrlInstall(BaseModel):
"""Request body for installing a plugin from a URL."""
url: str
MAX_UPLOAD_SIZE = 10 * 1024 * 1024 # 10 MB
@router.get("")
async def list_plugins(
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_permission("plugins:read")),
):
"""List all plugins with their current status (discovered, installed, active, inactive)."""
service = get_plugin_service()
plugins = await service.list_plugins(db)
return {"plugins": plugins, "total": len(plugins)}
@router.get("/manifest")
async def get_manifest_schema(
current_user: dict = Depends(require_permission("plugins:read")),
):
"""Get the plugin manifest schema documentation."""
service = get_plugin_service()
return service.get_manifest_schema()
@router.get("/updates")
async def check_plugin_updates(
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_permission("plugins:read")),
):
"""Check for available plugin updates.
Compares installed plugin versions with discovered plugin versions.
Returns a list of plugins where the discovered version is newer.
"""
from app.plugins.semver import SemVer
service = get_plugin_service()
plugins = await service.list_plugins(db)
updates: list[dict[str, Any]] = []
for plugin in plugins:
if not plugin.get("installed"):
continue
installed_version = plugin.get("version", "0.0.0")
# The discovered version is always the manifest version
discovered_version = plugin.get("version", "0.0.0")
# In a real marketplace scenario, we'd compare with a remote registry
# For now, we check if the manifest version differs from the DB version
# This is a placeholder for marketplace integration
try:
if SemVer.parse(discovered_version) > SemVer.parse(installed_version):
updates.append({
"name": plugin["name"],
"display_name": plugin.get("display_name", plugin["name"]),
"current_version": installed_version,
"available_version": discovered_version,
})
except ValueError:
pass # Skip if version is not valid SemVer
return {"updates": updates, "total": len(updates)}
@router.get("/active-manifests")
async def get_active_manifests(
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_permission("plugins:read")),
):
"""Get UI manifests for all active plugins.
Returns menu_items, page_routes, detail_tabs, settings_pages, and
dashboard_widgets contributed by each active plugin. Used by the
frontend PluginRegistry to dynamically register routes, sidebar items,
settings pages, and detail tabs.
"""
service = get_plugin_service()
manifests = await service.get_active_manifests(db)
return {"plugins": manifests, "total": len(manifests)}
@router.get("/{name}/config")
async def get_plugin_config(
name: str,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_permission("plugins:read")),
):
"""Get the configuration for a specific plugin.
Returns the plugin's config field as a JSON object.
"""
service = get_plugin_service()
try:
config = await service.get_plugin_config(db, name)
return {"name": name, "config": config}
except ValueError as exc:
raise HTTPException(
404, detail={"detail": str(exc), "code": "plugin_not_found"}
) from None
@router.patch("/{name}/config")
async def update_plugin_config(
name: str,
body: PluginConfigUpdate,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Update the configuration for a specific plugin.
Stores the config as a JSON string in the plugin's config field.
"""
import uuid as uuid_mod
service = get_plugin_service()
try:
result = await service.update_plugin_config(
db,
name,
config=body.config,
tenant_id=uuid_mod.UUID(current_user["tenant_id"]),
user_id=uuid_mod.UUID(current_user["user_id"]),
)
return result
except ValueError as exc:
raise HTTPException(
404, detail={"detail": str(exc), "code": "plugin_not_found"}
) from None
@router.post("/{name}/install")
async def install_plugin(
name: str,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Install a plugin by name. Runs migrations and creates DB record.
Idempotent: returns 200 if already installed.
"""
import uuid as uuid_mod
service = get_plugin_service()
try:
result = await service.install_plugin(
db,
name,
tenant_id=uuid_mod.UUID(current_user["tenant_id"]),
user_id=uuid_mod.UUID(current_user["user_id"]),
)
return result
except ValueError as exc:
if "not found" in str(exc).lower():
raise HTTPException(
404, detail={"detail": str(exc), "code": "plugin_not_found"}
) from None
raise HTTPException(400, detail={"detail": str(exc), "code": "plugin_error"}) from None
except MigrationValidationError as exc:
raise HTTPException(
422, detail={"detail": str(exc), "code": "migration_validation_error"}
) from None
@router.post("/{name}/activate")
async def activate_plugin(
name: str,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Activate a plugin by name. Registers event listeners and routes.
Idempotent: returns 200 if already active.
"""
import uuid as uuid_mod
service = get_plugin_service()
try:
result = await service.activate_plugin(
db,
name,
tenant_id=uuid_mod.UUID(current_user["tenant_id"]),
user_id=uuid_mod.UUID(current_user["user_id"]),
)
return result
except ValueError as exc:
if "not installed" in str(exc).lower():
raise HTTPException(
400, detail={"detail": str(exc), "code": "plugin_not_installed"}
) from None
if "not found" in str(exc).lower():
raise HTTPException(
404, detail={"detail": str(exc), "code": "plugin_not_found"}
) from None
raise HTTPException(400, detail={"detail": str(exc), "code": "plugin_error"}) from None
@router.post("/{name}/deactivate")
async def deactivate_plugin(
name: str,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Deactivate a plugin by name. Unregisters event listeners and routes.
Idempotent: returns 200 if already inactive.
"""
import uuid as uuid_mod
service = get_plugin_service()
try:
result = await service.deactivate_plugin(
db,
name,
tenant_id=uuid_mod.UUID(current_user["tenant_id"]),
user_id=uuid_mod.UUID(current_user["user_id"]),
)
return result
except ValueError as exc:
if "not found" in str(exc).lower():
raise HTTPException(
404, detail={"detail": str(exc), "code": "plugin_not_found"}
) from None
raise HTTPException(400, detail={"detail": str(exc), "code": "plugin_error"}) from None
@router.delete("/{name}")
async def uninstall_plugin(
name: str,
remove_data: bool = Query(False),
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Uninstall a plugin by name. Deactivates, drops tables, removes DB record.
Idempotent: returns 200 if already uninstalled.
"""
import uuid as uuid_mod
service = get_plugin_service()
try:
result = await service.uninstall_plugin(
db,
name,
remove_data=remove_data,
tenant_id=uuid_mod.UUID(current_user["tenant_id"]),
user_id=uuid_mod.UUID(current_user["user_id"]),
)
return result
except ValueError as exc:
if "not found" in str(exc).lower():
raise HTTPException(
404, detail={"detail": str(exc), "code": "plugin_not_found"}
) from None
raise HTTPException(400, detail={"detail": str(exc), "code": "plugin_error"}) from None
# ── Plugin Upload / URL Install ──────────────────────────────────────────────
def _validate_manifest_name(name: str) -> str:
"""Validate plugin name is alphanumeric with underscores only."""
if not re.match(r"^[a-zA-Z][a-zA-Z0-9_]*$", name):
raise ValueError(
f"Invalid plugin name '{name}': must start with a letter and contain only "
f"alphanumeric characters and underscores"
)
return name
def _check_dangerous_imports(source_code: str) -> list[str]:
"""Check plugin source for dangerous imports/patterns.
Returns a list of dangerous patterns found (empty if safe).
"""
dangerous_patterns = [
(r"\bos\.system\b", "os.system call"),
(r"\bsubprocess\.", "subprocess module"),
(r"\beval\s*\(", "eval() call"),
(r"\bexec\s*\(", "exec() call"),
(r"\b__import__\s*\(", "__import__() call"),
(r"\bcompile\s*\(", "compile() call"),
]
found: list[str] = []
for pattern, description in dangerous_patterns:
if re.search(pattern, source_code):
found.append(description)
return found
def _check_migration_sql(sql_content: str) -> list[str]:
"""Basic SQL validation for migration files.
Returns a list of issues found (empty if OK).
"""
issues: list[str] = []
# Check for basic SQL syntax issues
lines = sql_content.strip().split("\n")
for i, line in enumerate(lines, 1):
stripped = line.strip()
if not stripped or stripped.startswith("--"):
continue
# Check for unclosed parentheses
if stripped.count("(") != stripped.count(")"):
issues.append(f"Line {i}: unbalanced parentheses")
# Check for DROP TABLE (dangerous in migrations)
if re.search(r"\bDROP\s+TABLE\b", stripped, re.IGNORECASE):
issues.append(f"Line {i}: DROP TABLE is not allowed in plugin migrations")
return issues
def _find_plugin_class_in_module(module: Any) -> type[BasePlugin] | None:
"""Find a BasePlugin subclass in a module."""
for attr_name in dir(module):
attr = getattr(module, attr_name)
if (
isinstance(attr, type)
and issubclass(attr, BasePlugin)
and attr is not BasePlugin
):
return attr
return None
def _extract_plugin_from_zip(zip_path: str) -> tuple[Path, str, type[BasePlugin]]:
"""Extract a ZIP file and find the plugin class.
Returns (extract_dir, plugin_name, plugin_class).
"""
extract_dir = Path(tempfile.mkdtemp(prefix="plugin_upload_"))
try:
with zipfile.ZipFile(zip_path, "r") as zf:
# Validate ZIP
bad_files = [f for f in zf.namelist() if f.startswith("..") or f.startswith("/")]
if bad_files:
raise ValueError(f"ZIP contains files with unsafe paths: {bad_files}")
zf.extractall(extract_dir)
# Find plugin.py in the extracted contents
plugin_py_path: Path | None = None
for fpath in extract_dir.rglob("plugin.py"):
plugin_py_path = fpath
break
if plugin_py_path is None:
raise ValueError("ZIP does not contain a plugin.py file")
# Security: check source code BEFORE executing it
source_code = plugin_py_path.read_text(encoding="utf-8")
dangerous = _check_dangerous_imports(source_code)
if dangerous:
raise ValueError(
f"Plugin contains dangerous patterns: {', '.join(dangerous)}"
)
# Import the module dynamically (safe — source validated above)
spec = importlib.util.spec_from_file_location(
"uploaded_plugin", plugin_py_path
)
if spec is None or spec.loader is None:
raise ValueError("Could not load plugin.py module")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
# Find BasePlugin subclass
plugin_class = _find_plugin_class_in_module(module)
if plugin_class is None:
raise ValueError(
"plugin.py does not contain a BasePlugin subclass"
)
plugin_instance = plugin_class()
plugin_name = plugin_instance.name
# Validate manifest
manifest = plugin_instance.manifest
if not manifest.name or not manifest.version or not manifest.display_name:
raise ValueError(
"Plugin manifest must include name, version, and display_name"
)
# Validate plugin name
_validate_manifest_name(manifest.name)
# Check migration SQL files
migrations_dir = plugin_py_path.parent / "migrations"
if migrations_dir.exists():
for sql_file in sorted(migrations_dir.glob("*.sql")):
sql_content = sql_file.read_text(encoding="utf-8")
issues = _check_migration_sql(sql_content)
if issues:
raise ValueError(
f"Migration file {sql_file.name} has issues: {'; '.join(issues)}"
)
return extract_dir, plugin_name, plugin_class
except Exception:
# Clean up on failure
shutil.rmtree(extract_dir, ignore_errors=True)
raise
def _install_plugin_from_dir(
extract_dir: Path,
plugin_name: str,
plugin_class: type[BasePlugin],
) -> None:
"""Copy plugin directory to builtins and register it.
Copies the extracted plugin directory to app/plugins/builtins/{plugin_name}/.
"""
builtins_dir = Path(__file__).parent.parent / "plugins" / "builtins" / plugin_name
builtins_dir.mkdir(parents=True, exist_ok=True)
# Copy all files from extract_dir to builtins_dir
for item in extract_dir.iterdir():
dest = builtins_dir / item.name
if item.is_dir():
if dest.exists():
shutil.rmtree(dest)
shutil.copytree(item, dest)
else:
shutil.copy2(item, dest)
# Register the plugin in the registry
registry = get_plugin_service().registry
instance = plugin_class()
registry.register_plugin(instance)
logger.info(
"Installed plugin '%s' from uploaded ZIP to %s",
plugin_name,
builtins_dir,
)
@router.post("/upload")
async def upload_plugin(
file: UploadFile = File(...),
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Upload and install a plugin from a ZIP file.
DISABLED — Plugin upload is deactivated due to security vulnerabilities (RCE via exec_module before validation).
Will be re-enabled with signed plugin artifacts and sandboxed execution.
"""
raise HTTPException(
status_code=403,
detail={"detail": "Plugin upload is disabled. Use signed plugin artifacts from the allowlist.", "code": "upload_disabled"},
)
@router.post("/install-url")
async def install_plugin_from_url(
body: PluginUrlInstall,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Install a plugin from a URL (downloads ZIP and installs).
DISABLED — URL installation is deactivated due to SSRF and RCE vulnerabilities.
Will be re-enabled with signed plugin artifacts and allowlist.
"""
raise HTTPException(
status_code=403,
detail={"detail": "Plugin URL installation is disabled. Use signed plugin artifacts from the allowlist.", "code": "install_url_disabled"},
)
# ── Marketplace (Phase 5) ──────────────────────────────────────────────────────
class MarketplaceInstall(BaseModel):
"""Request body for installing a plugin from the marketplace."""
url: str
signature: str | None = None
public_key: str | None = None
activate: bool = False
@router.post("/install-marketplace")
async def install_from_marketplace(
body: MarketplaceInstall,
db: AsyncSession = Depends(get_db),
current_user: dict = Depends(require_admin),
):
"""Install a plugin from the marketplace.
1. Download ZIP from marketplace URL
2. Verify signature against allowlist (if provided)
3. Quarantine: validate manifest, check dangerous imports, validate SQL
4. Install (migrations + DB record)
5. Activate (optional)
DISABLED until marketplace is live — requires allowlist entry.
"""
raise HTTPException(
status_code=403,
detail={
"detail": "Marketplace installation is not yet available. Use built-in plugins.",
"code": "marketplace_not_available",
},
)