1055 lines
42 KiB
Python
1055 lines
42 KiB
Python
import asyncio
|
|
import html as html_mod
|
|
import json as _json
|
|
import logging
|
|
import os
|
|
import re
|
|
import secrets
|
|
import string
|
|
from contextlib import asynccontextmanager
|
|
from pathlib import Path
|
|
|
|
import mistune
|
|
from fastapi import Depends, FastAPI, HTTPException, Request, WebSocket
|
|
from fastapi.responses import FileResponse, HTMLResponse, JSONResponse, StreamingResponse
|
|
from fastapi.staticfiles import StaticFiles
|
|
from pydantic import BaseModel, Field
|
|
from starlette.middleware.base import BaseHTTPMiddleware
|
|
|
|
from backend.collab import authenticate_websocket, collab_manager
|
|
from backend.image_processor import preprocess_images
|
|
from backend.indexer import (
|
|
build_index,
|
|
find_file_in_index,
|
|
get_vault_data,
|
|
handle_file_move,
|
|
remove_single_file,
|
|
update_single_file,
|
|
)
|
|
from backend.openapi_docs import (
|
|
API_DESCRIPTION,
|
|
TAGS_METADATA,
|
|
enrich_openapi_schema,
|
|
render_api_landing,
|
|
)
|
|
from backend.search import (
|
|
init_inverted_index,
|
|
)
|
|
from backend.semantic_search import init_semantic_index
|
|
from backend.services.backups import get_backup_dir as service_get_backup_dir
|
|
from backend.services.errors import ServiceError
|
|
from backend.services.sanitizer import sanitize_html
|
|
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="%(asctime)s [%(name)s] %(levelname)s: %(message)s",
|
|
)
|
|
logger = logging.getLogger("obsigate")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Pydantic response models : voir backend.schemas (vaults/history : #85 T8)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# Filesystem mutation + search / suggest / graph models : voir backend.schemas (#85 T5, T6b)
|
|
|
|
class BackupEntry(BaseModel):
|
|
"""A single backup version of a file."""
|
|
timestamp: int = Field(description="Unix timestamp of when the backup was created")
|
|
datetime: str = Field(description="ISO 8601 datetime string")
|
|
size: int = Field(description="File size in bytes")
|
|
filename: str = Field(description="Backup filename on disk")
|
|
|
|
|
|
class BackupListResponse(BaseModel):
|
|
"""Response listing all available backups for a file."""
|
|
vault: str = Field(description="Vault name")
|
|
path: str = Field(description="Relative file path")
|
|
backups: list[BackupEntry] = Field(description="Available backups, newest first")
|
|
|
|
|
|
class DiffRequest(BaseModel):
|
|
"""Request parameters for generating a diff."""
|
|
version: int = Field(description="Timestamp of the backup version to compare")
|
|
compare_with: int | None = Field(default=None, description="Timestamp of another backup version. If omitted, compares with the current file.")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# SSE Manager — voir backend.sse (ROADMAP #85 T4, instance partagée)
|
|
# ---------------------------------------------------------------------------
|
|
# ---------------------------------------------------------------------------
|
|
|
|
from backend.search_executor import (
|
|
get_search_executor,
|
|
init_search_executor,
|
|
shutdown_search_executor,
|
|
)
|
|
from backend.sse import sse_manager
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Application lifespan (replaces deprecated on_event)
|
|
# ---------------------------------------------------------------------------
|
|
from backend.watcher import VaultWatcher
|
|
from backend.watcher_state import get_watcher, set_watcher
|
|
|
|
# File watcher : handle partagé via backend.watcher_state (ROADMAP #85 T8).
|
|
|
|
|
|
async def _on_vault_change(events: list):
|
|
"""Callback invoked by VaultWatcher when files change in watched vaults.
|
|
|
|
Processes each event (create/modify/delete/move) and updates the index
|
|
incrementally, then broadcasts SSE notifications.
|
|
"""
|
|
updated_vaults = set()
|
|
changes = []
|
|
|
|
for event in events:
|
|
vault_name = event["vault"]
|
|
event_type = event["type"]
|
|
src = event["src"]
|
|
dest = event.get("dest")
|
|
|
|
try:
|
|
if event_type in ("created", "modified"):
|
|
result = await update_single_file(vault_name, src)
|
|
if result:
|
|
changes.append({"action": "updated", "vault": vault_name, "path": result["path"]})
|
|
updated_vaults.add(vault_name)
|
|
|
|
elif event_type == "deleted":
|
|
result = await remove_single_file(vault_name, src)
|
|
if result:
|
|
changes.append({"action": "deleted", "vault": vault_name, "path": result["path"]})
|
|
updated_vaults.add(vault_name)
|
|
|
|
elif event_type == "moved":
|
|
result = await handle_file_move(vault_name, src, dest)
|
|
if result:
|
|
changes.append({"action": "moved", "vault": vault_name, "path": result["path"]})
|
|
updated_vaults.add(vault_name)
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error processing {event_type} event for {src}: {e}")
|
|
|
|
if changes:
|
|
await sse_manager.broadcast("index_updated", {
|
|
"vaults": list(updated_vaults),
|
|
"changes": changes,
|
|
"total_changes": len(changes),
|
|
})
|
|
logger.info(f"Hot-reload: {len(changes)} change(s) in {list(updated_vaults)}")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Authentication bootstrap
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def bootstrap_admin():
|
|
"""Create the initial admin account if no users exist.
|
|
|
|
Reads OBSIGATE_ADMIN_USER and OBSIGATE_ADMIN_PASSWORD from environment.
|
|
If no password is set, generates a random one and logs it ONCE.
|
|
Only runs when auth is enabled and no users.json exists yet.
|
|
"""
|
|
from backend.auth.middleware import is_auth_enabled
|
|
from backend.auth.user_store import create_user, has_users
|
|
|
|
if not is_auth_enabled():
|
|
return
|
|
|
|
if has_users():
|
|
return # Users already exist, skip
|
|
|
|
admin_user = os.environ.get("OBSIGATE_ADMIN_USER", "admin")
|
|
admin_pass = os.environ.get("OBSIGATE_ADMIN_PASSWORD", "")
|
|
|
|
if not admin_pass:
|
|
# Generate a random password and display it ONCE in logs
|
|
admin_pass = "".join(
|
|
secrets.choice(string.ascii_letters + string.digits)
|
|
for _ in range(16)
|
|
)
|
|
logger.warning("=" * 60)
|
|
logger.warning("PREMIER DÉMARRAGE — Compte admin créé automatiquement")
|
|
logger.warning(f" Utilisateur : {admin_user}")
|
|
logger.warning(f" Mot de passe : {admin_pass}")
|
|
logger.warning("CHANGEZ CE MOT DE PASSE dès la première connexion !")
|
|
logger.warning("=" * 60)
|
|
|
|
try:
|
|
create_user(admin_user, admin_pass, role="admin", vaults=["*"])
|
|
logger.info(f"Admin '{admin_user}' créé avec succès")
|
|
except PermissionError as e:
|
|
logger.critical("=" * 60)
|
|
logger.critical("DÉMARRAGE IMPOSSIBLE : Erreur de permission sur le dossier 'data'")
|
|
logger.critical("L'indexation et l'authentification ne peuvent pas fonctionner.")
|
|
logger.critical("FIX : Vérifiez les droits du volume /app/data sur l'hôte.")
|
|
logger.critical("Exemple : sudo chown -R 1000:1000 /DOCKER_CONFIG/ObsiGate/data")
|
|
logger.critical("=" * 60)
|
|
raise e
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Security headers middleware
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class SecurityHeadersMiddleware(BaseHTTPMiddleware):
|
|
"""Add security headers to all HTTP responses."""
|
|
|
|
async def dispatch(self, request, call_next):
|
|
response = await call_next(request)
|
|
response.headers["X-Content-Type-Options"] = "nosniff"
|
|
response.headers["X-Frame-Options"] = "SAMEORIGIN"
|
|
response.headers["X-XSS-Protection"] = "1; mode=block"
|
|
response.headers["Referrer-Policy"] = "strict-origin-when-cross-origin"
|
|
# A route may set a stricter per-response policy (e.g. ``sandbox`` for
|
|
# standalone SVG, #108-B3); keep it instead of overwriting it.
|
|
if "Content-Security-Policy" not in response.headers:
|
|
response.headers["Content-Security-Policy"] = (
|
|
"default-src 'self'; "
|
|
"script-src 'self' 'unsafe-inline' blob: https://cdnjs.cloudflare.com https://unpkg.com https://esm.sh https://cdn.jsdelivr.net https://static.cloudflareinsights.com; "
|
|
"style-src 'self' 'unsafe-inline' https://cdnjs.cloudflare.com https://fonts.googleapis.com https://cdn.jsdelivr.net https://esm.sh; "
|
|
"img-src 'self' data: blob:; "
|
|
"connect-src 'self' blob: https://esm.sh https://unpkg.com https://cdnjs.cloudflare.com https://fonts.googleapis.com https://fonts.gstatic.com https://cdn.jsdelivr.net; "
|
|
"font-src 'self' data: https://fonts.gstatic.com https://esm.sh; "
|
|
"worker-src 'self' blob:; "
|
|
"frame-src 'self' blob:; "
|
|
"object-src 'none'; "
|
|
"base-uri 'self'; "
|
|
"form-action 'self'; "
|
|
"frame-ancestors 'self';"
|
|
)
|
|
# Static assets are NOT content-hashed, so they must revalidate:
|
|
# ``immutable``/long max-age made Cloudflare and mobile browsers serve
|
|
# a stale build for a year (the service worker cache compounded it).
|
|
# ``no-cache`` keeps caching but forces revalidation (ETag/Last-Modified).
|
|
if request.url.path.startswith("/static/"):
|
|
response.headers["Cache-Control"] = "no-cache"
|
|
return response
|
|
|
|
|
|
def _guard_insecure_auth() -> None:
|
|
"""Warn or refuse to start when authentication is disabled (BUG-037).
|
|
|
|
With ``OBSIGATE_AUTH_ENABLED=false`` every request is served as an
|
|
anonymous admin. That is convenient for local use but dangerous when the
|
|
process is reachable from a network. Binding to a non-loopback host
|
|
without the explicit ``OBSIGATE_ALLOW_INSECURE=true`` opt-in is refused.
|
|
"""
|
|
from backend.auth.middleware import (
|
|
bind_host_from_argv,
|
|
is_auth_enabled,
|
|
is_insecure_mode_allowed,
|
|
is_loopback_host,
|
|
)
|
|
|
|
if is_auth_enabled():
|
|
return
|
|
|
|
if is_insecure_mode_allowed():
|
|
logger.warning(
|
|
"Authentication is DISABLED and OBSIGATE_ALLOW_INSECURE=true: every request "
|
|
"is treated as an anonymous administrator. Do not expose this instance."
|
|
)
|
|
return
|
|
|
|
host = bind_host_from_argv()
|
|
if not is_loopback_host(host):
|
|
raise RuntimeError(
|
|
"Refusing to start: authentication is disabled (OBSIGATE_AUTH_ENABLED=false) "
|
|
f"while binding to a non-loopback address ('{host}'). This would expose an "
|
|
"unauthenticated instance with admin access. Enable authentication, or set "
|
|
"OBSIGATE_ALLOW_INSECURE=true if you really know what you are doing."
|
|
)
|
|
|
|
logger.warning(
|
|
"Authentication is DISABLED (OBSIGATE_AUTH_ENABLED=false): every request is "
|
|
"treated as an anonymous administrator. This is only safe on a trusted, "
|
|
"loopback-only deployment."
|
|
)
|
|
|
|
|
|
@asynccontextmanager
|
|
async def lifespan(app: FastAPI):
|
|
"""Application lifespan: build index on startup, cleanup on shutdown."""
|
|
# Thread pool for offloading CPU-bound search from the event loop.
|
|
# Sized to 2 workers so concurrent searches don't starve other requests.
|
|
init_search_executor()
|
|
|
|
# BUG-037: refuse to expose an unauthenticated instance on a public bind.
|
|
_guard_insecure_auth()
|
|
|
|
# Bootstrap admin account if needed
|
|
bootstrap_admin()
|
|
|
|
logger.info("ObsiGate starting — building index in background...")
|
|
|
|
async def _progress_cb(event_type: str, data: dict):
|
|
await sse_manager.broadcast("index_" + event_type, data)
|
|
|
|
async def _background_startup():
|
|
logger.info("Background indexing started")
|
|
await build_index(_progress_cb)
|
|
|
|
# Build inverted index in a thread pool to avoid blocking the event loop.
|
|
# The inverted index rebuild is CPU-bound (tokenization, indexing) and
|
|
# would freeze HTTP responses if run in the async event loop.
|
|
loop = asyncio.get_running_loop()
|
|
await loop.run_in_executor(get_search_executor(), init_inverted_index)
|
|
# Build the semantic (embedding) index in the same background thread pool.
|
|
await loop.run_in_executor(get_search_executor(), init_semantic_index)
|
|
|
|
# BUG-040: extract the PDF text deferred during the scan now that the
|
|
# index and inverted index are queryable (keeps startup non-blocking).
|
|
from backend.indexer import enrich_pdf_texts
|
|
await enrich_pdf_texts()
|
|
|
|
# Scan for plugins in all vaults
|
|
logger.info("Scanning for plugins...")
|
|
from backend.indexer import vault_config
|
|
from backend.plugins import get_plugin_registry
|
|
registry = get_plugin_registry()
|
|
for vault_name, cfg in vault_config.items():
|
|
vault_path = cfg.get("path")
|
|
if vault_path:
|
|
try:
|
|
plugins = registry.scan_vault(vault_name, vault_path)
|
|
logger.info(f"Vault '{vault_name}': found {len(plugins)} plugin(s)")
|
|
from backend.plugins import emit_vault_mounted
|
|
emit_vault_mounted(vault_name, vault_path)
|
|
except Exception as e:
|
|
logger.warning(f"Plugin scan failed for vault '{vault_name}': {e}")
|
|
|
|
# Start file watcher (handle partagé : voir backend.watcher_state)
|
|
config = _load_config()
|
|
watcher_enabled = config.get("watcher_enabled", True)
|
|
if watcher_enabled:
|
|
use_polling = config.get("watcher_use_polling", False)
|
|
polling_interval = config.get("watcher_polling_interval", 5.0)
|
|
debounce = config.get("watcher_debounce", 2.0)
|
|
watcher = VaultWatcher(
|
|
on_file_change=_on_vault_change,
|
|
debounce_seconds=debounce,
|
|
use_polling=use_polling,
|
|
polling_interval=polling_interval,
|
|
)
|
|
from backend.indexer import vault_config
|
|
vaults_to_watch = {name: cfg["path"] for name, cfg in vault_config.items()}
|
|
await watcher.start(vaults_to_watch)
|
|
set_watcher(watcher)
|
|
logger.info("File watcher started in background.")
|
|
else:
|
|
logger.info("File watcher disabled by configuration.")
|
|
|
|
logger.info("Background startup complete.")
|
|
|
|
asyncio.create_task(_background_startup())
|
|
|
|
logger.info("ObsiGate ready (listening for requests while indexing).")
|
|
yield
|
|
|
|
# Shutdown
|
|
await collab_manager.stop()
|
|
watcher = get_watcher()
|
|
if watcher:
|
|
await watcher.stop()
|
|
set_watcher(None)
|
|
shutdown_search_executor()
|
|
|
|
|
|
from backend.version import get_version
|
|
|
|
app = FastAPI(
|
|
title="ObsiGate API",
|
|
version=get_version(),
|
|
lifespan=lifespan,
|
|
description=API_DESCRIPTION.strip(),
|
|
openapi_tags=TAGS_METADATA,
|
|
docs_url="/docs",
|
|
redoc_url="/redoc",
|
|
openapi_url="/openapi.json",
|
|
contact={"name": "ObsiGate", "url": "https://git.dracodev.net/Projets/ObsiGate"},
|
|
license_info={"name": "MIT"},
|
|
)
|
|
|
|
# Enrich the auto-generated OpenAPI 3.1 schema (#72): tags per category,
|
|
# examples, security schemes and documented error responses.
|
|
_original_openapi = app.openapi
|
|
|
|
|
|
def _custom_openapi():
|
|
if app.openapi_schema:
|
|
return app.openapi_schema
|
|
schema = _original_openapi()
|
|
app.openapi_schema = enrich_openapi_schema(schema)
|
|
return app.openapi_schema
|
|
|
|
|
|
app.openapi = _custom_openapi # type: ignore[method-assign]
|
|
|
|
|
|
@app.exception_handler(ServiceError)
|
|
async def _service_error_handler(request: Request, exc: ServiceError):
|
|
"""Map shared-layer domain errors to HTTP responses (``{"detail": ...}``)."""
|
|
return JSONResponse(status_code=exc.status, content={"detail": exc.message})
|
|
|
|
# GZip compression — reduces bandwidth by ~70% for text responses
|
|
# Custom wrapper: skip compression for SSE streams (/api/events)
|
|
from fastapi.middleware.gzip import GZipMiddleware
|
|
from starlette.types import Receive, Scope, Send
|
|
|
|
|
|
class SSESafeGZipMiddleware(GZipMiddleware):
|
|
"""GZip middleware that skips SSE (Server-Sent Events) streams.
|
|
|
|
GZip buffering breaks incremental streaming required by SSE.
|
|
We detect SSE endpoints by path and bypass compression entirely.
|
|
"""
|
|
# SSE endpoints that must not be buffered by GZip.
|
|
_SSE_PATHS = (
|
|
"/api/events",
|
|
"/api/admin/stream",
|
|
"/api/ai/bookslm/chat",
|
|
"/api/ai/bookslm/agent",
|
|
"/mcp",
|
|
)
|
|
|
|
async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None:
|
|
if scope["type"] == "http" and scope.get("path") in self._SSE_PATHS:
|
|
# Bypass GZip: passthrough directly to the inner app
|
|
await self.app(scope, receive, send)
|
|
else:
|
|
await super().__call__(scope, receive, send)
|
|
|
|
app.add_middleware(SSESafeGZipMiddleware, minimum_size=1000)
|
|
|
|
# Security headers on all responses
|
|
app.add_middleware(SecurityHeadersMiddleware)
|
|
|
|
# Auth router
|
|
# Multi-format export (HTML / MD bundle / ePub) — voir backend.routers.files_media (#85 T6c).
|
|
from backend.ai_routes import router as ai_router
|
|
from backend.auth.middleware import (
|
|
check_vault_access,
|
|
require_admin,
|
|
require_auth,
|
|
)
|
|
from backend.auth.router import router as auth_router
|
|
from backend.bookslm_routes import router as bookslm_router
|
|
from backend.routers.backups import router as backups_router
|
|
from backend.routers.config import _load_config
|
|
from backend.routers.config import router as config_router
|
|
from backend.routers.conflicts import router as conflicts_router
|
|
from backend.routers.files_media import router as files_media_router
|
|
from backend.routers.files_read import router as files_read_router
|
|
from backend.routers.files_write import router as files_write_router
|
|
from backend.routers.health import router as health_router
|
|
from backend.routers.history import router as history_router
|
|
from backend.routers.search import router as search_router
|
|
from backend.routers.sharing import router as sharing_router
|
|
from backend.routers.vaults import router as vaults_router
|
|
from backend.routers.webhooks import router as webhooks_router
|
|
from backend.secret_redactor import redact_file_content
|
|
from backend.skills_routes import router as skills_router
|
|
|
|
app.include_router(auth_router)
|
|
app.include_router(ai_router)
|
|
app.include_router(bookslm_router)
|
|
app.include_router(skills_router)
|
|
app.include_router(health_router) # ROADMAP #85 T1 — System / health
|
|
app.include_router(history_router) # ROADMAP #85 T8 — History
|
|
app.include_router(search_router) # ROADMAP #85 T5 — Search
|
|
app.include_router(backups_router) # ROADMAP #85 T4 — Backups
|
|
app.include_router(conflicts_router) # ROADMAP #85 T8 — Conflicts
|
|
app.include_router(config_router) # ROADMAP #85 T7 — Config
|
|
app.include_router(files_read_router) # ROADMAP #85 T6a — Files read
|
|
app.include_router(files_media_router) # ROADMAP #85 T6c — Media/export
|
|
app.include_router(files_write_router) # ROADMAP #85 T6b — Files write
|
|
app.include_router(webhooks_router) # ROADMAP #85 T2 — Webhooks
|
|
app.include_router(sharing_router) # ROADMAP #85 T3 — Sharing
|
|
app.include_router(vaults_router) # ROADMAP #85 T8 — Vaults
|
|
|
|
# Admin Dashboard endpoints (system stats, audit logs, backups, stream)
|
|
try:
|
|
from backend.admin import router as admin_router
|
|
app.include_router(admin_router)
|
|
logger.info("Admin dashboard router mounted at /api/admin/*")
|
|
except ImportError as e:
|
|
logger.warning(f"Could not load admin dashboard router: {e}")
|
|
|
|
# Push Notifications endpoints (Web Push API + VAPID)
|
|
try:
|
|
from backend.push import router as push_router
|
|
app.include_router(push_router)
|
|
logger.info("Push notifications router mounted at /api/push/*")
|
|
except ImportError as e:
|
|
logger.warning(f"Could not load push notifications router: {e}")
|
|
|
|
# Plugins system endpoints
|
|
try:
|
|
from backend.plugins import router as plugins_router
|
|
app.include_router(plugins_router)
|
|
logger.info("Plugins router mounted at /api/plugins/*")
|
|
except ImportError as e:
|
|
logger.warning(f"Could not load plugins router: {e}")
|
|
|
|
# MCP server (Streamable HTTP) for external clients (#79 phase E)
|
|
try:
|
|
from backend.mcp.server import McpMount, mcp_app
|
|
app.router.routes.append(McpMount(mcp_app))
|
|
logger.info("MCP server mounted at /mcp")
|
|
except Exception as e: # pragma: no cover - optional dependency
|
|
logger.warning(f"Could not mount MCP server: {e}")
|
|
|
|
# Resolve frontend path relative to this file
|
|
FRONTEND_DIR = Path(__file__).resolve().parent.parent / "frontend"
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# API documentation landing page (#72)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
@app.get("/api", include_in_schema=False, response_class=HTMLResponse)
|
|
@app.get("/api/", include_in_schema=False, response_class=HTMLResponse)
|
|
async def api_docs_landing():
|
|
"""Human-friendly API documentation landing page (links to /docs, /redoc)."""
|
|
return HTMLResponse(render_api_landing(get_version()))
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Path safety helper : voir backend.routers.helpers (#85 T6a, T6c)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _resolve_safe_path(vault_root: Path, relative_path: str | None) -> Path:
|
|
"""Resolve a relative path safely within the vault root.
|
|
|
|
Thin wrapper around the shared :func:`backend.services.paths.resolve_safe_path`
|
|
(single implementation used by both routes and tools). The raised
|
|
:class:`ServiceError` is mapped to an ``HTTPException`` response by the
|
|
global exception handler in this module.
|
|
|
|
Args:
|
|
vault_root: The vault's root directory (absolute).
|
|
relative_path: The user-supplied relative path.
|
|
|
|
Returns:
|
|
Resolved absolute ``Path``.
|
|
"""
|
|
from backend.services.paths import resolve_safe_path as _service_resolve
|
|
|
|
return _service_resolve(vault_root, relative_path)
|
|
|
|
|
|
def _backup_file(file_path: Path, vault_name: str, relative_path: str):
|
|
"""Create a timestamped backup of a file before modification.
|
|
|
|
Thin wrapper around :func:`backend.services.backups.create_backup`
|
|
(single implementation used by both routes and tools). Backups are stored
|
|
in ``{backup_root}/{vault}/{relative_path}.{timestamp}.bak``; the operation
|
|
is best-effort and never blocks the caller.
|
|
"""
|
|
from backend.services.backups import create_backup
|
|
|
|
create_backup(file_path, vault_name, relative_path)
|
|
|
|
|
|
def _check_vault_writable(vault_root: Path) -> bool:
|
|
"""Check if a vault is writable (not mounted read-only).
|
|
|
|
Args:
|
|
vault_root: The vault's root directory (absolute).
|
|
|
|
Returns:
|
|
True if the vault is writable, False otherwise.
|
|
"""
|
|
return os.access(vault_root, os.W_OK)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Markdown rendering helpers (singleton renderer)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
import unicodedata
|
|
|
|
|
|
def _heading_slugify(text: str) -> str:
|
|
"""Generate a URL-safe slug from heading text.
|
|
|
|
Matches the JavaScript slugify algorithm exactly using
|
|
Unicode-aware character classification:
|
|
1. Strip HTML tags (e.g. wikilink spans rendered inside headings)
|
|
2. Decode HTML entities (e.g. ``&`` → ``&``)
|
|
3. Lowercase
|
|
4. NFD normalize + strip combining marks
|
|
5. Keep only Unicode letters, numbers, spaces, hyphens
|
|
6. Replace spaces with hyphens, collapse multiple hyphens
|
|
|
|
Args:
|
|
text: The heading text content (may contain inline HTML).
|
|
|
|
Returns:
|
|
A URL-safe slug string.
|
|
"""
|
|
# Strip any inline HTML so it does not pollute the slug
|
|
text = re.sub(r"<[^>]+>", "", text)
|
|
# Decode HTML entities so & becomes & before slugification
|
|
text = html_mod.unescape(text)
|
|
text = text.lower()
|
|
text = unicodedata.normalize("NFD", text)
|
|
text = "".join(ch for ch in text if not unicodedata.combining(ch))
|
|
# Unicode-aware: keep letters (L*), numbers (N*), spaces, and hyphens
|
|
cleaned = []
|
|
for ch in text:
|
|
cat = unicodedata.category(ch)
|
|
if cat.startswith('L') or cat.startswith('N') or ch in (' ', '-'):
|
|
cleaned.append(ch)
|
|
text = "".join(cleaned)
|
|
text = re.sub(r"\s+", "-", text)
|
|
text = re.sub(r"-+", "-", text)
|
|
result = text.strip("-")
|
|
return result if result else "heading"
|
|
|
|
|
|
def _add_heading_ids(html: str) -> str:
|
|
"""Post-process rendered HTML to add IDs to heading tags.
|
|
|
|
Adds an ``id`` attribute to every ``<h1>`` through ``<h6>`` tag
|
|
using a slug generated from the heading's text content.
|
|
Duplicate slugs get a ``-2``, ``-3``, etc. suffix.
|
|
|
|
Args:
|
|
html: Rendered HTML string.
|
|
|
|
Returns:
|
|
HTML with heading IDs injected.
|
|
"""
|
|
used_ids: dict[str, int] = {}
|
|
|
|
def _replace_heading(match):
|
|
tag = match.group(1)
|
|
content = match.group(2)
|
|
slug = _heading_slugify(content)
|
|
count = used_ids.get(slug, 0)
|
|
used_ids[slug] = count + 1
|
|
if count > 0:
|
|
slug = f"{slug}-{count + 1}"
|
|
return f'<{tag} id="{slug}">{content}</{tag}>'
|
|
|
|
# Match h1-h6 tags with text content (no existing id attribute)
|
|
return re.sub(
|
|
r'<(h[1-6])>([^<]*(?:<(?!/?h[1-6])[^<]*)*)</h[1-6]>',
|
|
_replace_heading,
|
|
html,
|
|
)
|
|
|
|
|
|
# Cached mistune renderer — avoids re-creating on every request
|
|
_markdown_renderer = mistune.create_markdown(
|
|
escape=False,
|
|
plugins=["table", "strikethrough", "footnotes", "task_lists"],
|
|
)
|
|
|
|
|
|
def _convert_wikilinks(content: str, current_vault: str) -> str:
|
|
"""Convert ``[[wikilinks]]`` and ``[[target|display]]`` to clickable HTML.
|
|
|
|
Supports:
|
|
- Internal file links: ``[[My Note]]`` / ``[[My Note|display]]``
|
|
- Same-document anchors: ``[[#Heading]]`` / ``[[#Heading|display]]``
|
|
|
|
Resolved file links get a ``data-vault`` / ``data-path`` attribute pair.
|
|
Anchor links target the slugified heading ID in the current document.
|
|
Unresolved links are rendered as ``<span class="wikilink-missing">``.
|
|
|
|
Args:
|
|
content: Markdown string potentially containing wikilinks.
|
|
current_vault: Active vault name for resolution priority.
|
|
|
|
Returns:
|
|
Markdown string with wikilinks replaced by HTML anchors.
|
|
"""
|
|
def _replace(match):
|
|
target = match.group(1).strip()
|
|
display = match.group(2).strip() if match.group(2) else target
|
|
|
|
# Same-document anchor link: [[#Heading|display]]
|
|
if target.startswith("#"):
|
|
anchor_text = target[1:].strip()
|
|
anchor_slug = _heading_slugify(anchor_text)
|
|
link_display = display if display != target else anchor_text
|
|
return f'<a class="wikilink-anchor" href="#{anchor_slug}">{link_display}</a>'
|
|
|
|
found = find_file_in_index(target, current_vault)
|
|
if found:
|
|
return (
|
|
f'<a class="wikilink" href="#" '
|
|
f'data-vault="{found["vault"]}" '
|
|
f'data-path="{found["path"]}">{display}</a>'
|
|
)
|
|
return f'<span class="wikilink-missing">{display}</span>'
|
|
|
|
pattern = r'\[\[([^\]|]+)(?:\|([^\]]+))?\]\]'
|
|
return re.sub(pattern, _replace, content)
|
|
|
|
|
|
def _normalize_line_breaks(text: str) -> str:
|
|
"""Convert single newlines to hard breaks (matching Obsidian default behavior).
|
|
|
|
In standard Markdown, a single ``\\n`` is a "soft break" — it renders as a space,
|
|
not a visible line break. Obsidian defaults to treating single newlines as hard
|
|
breaks (equivalent to ``<br>``). This function pre-processes the Markdown source
|
|
so that mistune renders standalone lines on separate rows, while still honouring
|
|
blank lines as paragraph separators.
|
|
|
|
Fenced code blocks (`` ``` ``) are left untouched so their internal newlines are
|
|
preserved verbatim.
|
|
"""
|
|
parts = re.split(r"(```[\s\S]*?```)", text)
|
|
for i, part in enumerate(parts):
|
|
if part.startswith("```"):
|
|
continue # Protect fenced code blocks
|
|
# Single \n (not preceded or followed by another \n) → two spaces + \n
|
|
parts[i] = re.sub(r"(?<!\n)\n(?!\n)", " \n", part)
|
|
return "".join(parts)
|
|
|
|
|
|
def _render_markdown(raw_md: str, vault_name: str, current_file_path: Path | None = None) -> str:
|
|
"""Render a markdown string to HTML with wikilink and image support.
|
|
|
|
Uses the cached singleton mistune renderer for performance.
|
|
|
|
Args:
|
|
raw_md: Raw markdown text (frontmatter already stripped).
|
|
vault_name: Current vault for wikilink resolution context.
|
|
current_file_path: Absolute path to the current markdown file.
|
|
|
|
Returns:
|
|
HTML string.
|
|
"""
|
|
# Get vault data for image resolution
|
|
vault_data = get_vault_data(vault_name)
|
|
vault_root = Path(vault_data["path"]) if vault_data else None
|
|
attachments_path = vault_data.get("config", {}).get("attachmentsPath") if vault_data else None
|
|
|
|
# Redact secrets before rendering (P0 security)
|
|
raw_md = redact_file_content(raw_md, str(current_file_path) if current_file_path else "")
|
|
|
|
# Preprocess images first
|
|
if vault_root:
|
|
raw_md = preprocess_images(raw_md, vault_name, vault_root, current_file_path, attachments_path)
|
|
|
|
# Convert wikilinks
|
|
converted = _convert_wikilinks(raw_md, vault_name)
|
|
|
|
# Normalize line breaks to match Obsidian behavior (single \n → hard break)
|
|
converted = _normalize_line_breaks(converted)
|
|
|
|
rendered = _markdown_renderer(converted)
|
|
|
|
# Add heading IDs for TOC navigation
|
|
rendered = _add_heading_ids(rendered)
|
|
|
|
# Sanitize: raw HTML in vault content must never reach the DOM (BUG-021).
|
|
rendered = sanitize_html(rendered)
|
|
|
|
return rendered
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# API Endpoints — System / health : voir backend.routers.health (#85 T1)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Vaults : voir backend.routers.vaults (#85 T8)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# History (recent / bookmarks / saved-searches) : voir backend.routers.history (#85 T8)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# File browse & read endpoints : voir backend.routers.files_read (#85 T6a)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# File PDF / export / guide / media endpoints : voir backend.routers.files_media (#85 T6c)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# File & directory mutations : voir backend.routers.files_write (#85 T6b)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# File creation and rename endpoints : voir backend.routers.files_write (#85 T6b)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# File creation and rename endpoints : voir backend.routers.files_write (#85 T6b)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Backup & Diff endpoints
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _get_backup_dir(vault_name: str, relative_path: str) -> Path:
|
|
"""Return the directory where backups for a specific file are stored.
|
|
|
|
Thin wrapper around :func:`backend.services.backups.get_backup_dir`
|
|
(single implementation used by both routes and tools).
|
|
"""
|
|
return service_get_backup_dir(vault_name, relative_path)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# File-level backup endpoints : voir backend.routers.backups (#85 T4)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# File backlinks + view endpoints : voir backend.routers.files_read (#85 T6a)
|
|
|
|
|
|
# Range helper : voir backend.routers.helpers.stream_file_with_range (#85 T6c)
|
|
|
|
# PDF stream/info endpoints : voir backend.routers.files_media (#85 T6c)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Search / suggest / graph / index-reload : voir backend.routers.search (#85 T5)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# SSE endpoint — Server-Sent Events stream
|
|
# ---------------------------------------------------------------------------
|
|
|
|
@app.get(
|
|
"/api/events",
|
|
response_class=StreamingResponse,
|
|
responses={200: {"content": {"text/event-stream": {}}, "description": "Server-Sent Events stream"}},
|
|
)
|
|
async def api_events(current_user=Depends(require_auth)):
|
|
"""SSE stream for real-time index update notifications.
|
|
|
|
Sends keepalive comments every 30s. Events:
|
|
- ``index_updated``: partial index change (file create/modify/delete/move)
|
|
- ``index_reloaded``: full re-index completed
|
|
- ``vault_added``: new vault added dynamically
|
|
- ``vault_removed``: vault removed dynamically
|
|
"""
|
|
queue = await sse_manager.connect()
|
|
|
|
async def event_generator():
|
|
try:
|
|
# Send initial connection event
|
|
yield f"event: connected\ndata: {_json.dumps({'sse_clients': sse_manager.client_count})}\n\n"
|
|
while True:
|
|
try:
|
|
msg = await asyncio.wait_for(queue.get(), timeout=30.0)
|
|
yield f"event: {msg['event']}\ndata: {msg['data']}\n\n"
|
|
except asyncio.TimeoutError:
|
|
# Keepalive comment
|
|
yield ": keepalive\n\n"
|
|
except asyncio.CancelledError:
|
|
break
|
|
finally:
|
|
sse_manager.disconnect(queue)
|
|
|
|
return StreamingResponse(
|
|
event_generator(),
|
|
media_type="text/event-stream",
|
|
headers={
|
|
"Cache-Control": "no-cache",
|
|
"Connection": "keep-alive",
|
|
"X-Accel-Buffering": "no",
|
|
},
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Dynamic vault management endpoints : voir backend.routers.vaults (#85 T8)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Image / media / attachments / vault-settings : voir backend.routers.files_media (#85 T6c)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Vault Settings API — Display preferences : voir backend.routers.files_media (#85 T6c)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Backup Management API
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Global backup endpoints : voir backend.routers.backups (#85 T4)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Configuration API : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AI API Keys : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tool & connected-source keys (#103) : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AI Models : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# AI Models : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Diagnostics API : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Dashboard endpoint (aggregated stats) : voir backend.routers.config (#85 T7)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Webhook CRUD endpoints : voir backend.routers.webhooks (#85 T2)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Share (public document) endpoints : voir backend.routers.sharing (#85 T3)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Syncthing conflict endpoints : voir backend.routers.conflicts (#85 T8)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Real-time collaboration — WebSocket endpoint (ROADMAP #62)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
@app.websocket("/ws/collab/{vault_name}/{path:path}")
|
|
async def collab_websocket(websocket: WebSocket, vault_name: str, path: str):
|
|
"""Real-time collaborative editing over WebSocket (ROADMAP #62).
|
|
|
|
One *room* is created per ``vault::path``; all clients editing the same
|
|
file share Yjs/CRDT updates, awareness (cursors/selection) and a debounced
|
|
server-side persistence of the markdown content.
|
|
|
|
Authentication is performed manually (FastAPI ``Depends`` do not run for
|
|
WebSocket routes) and vault access is enforced per connection.
|
|
"""
|
|
from backend.services.errors import ServiceError
|
|
from backend.services.vaults import get_vault_root
|
|
|
|
user = authenticate_websocket(websocket)
|
|
if user is None:
|
|
await websocket.close(code=4401)
|
|
return
|
|
|
|
if not check_vault_access(vault_name, user):
|
|
await websocket.close(code=4403)
|
|
return
|
|
|
|
try:
|
|
vault_root = get_vault_root(vault_name)
|
|
file_path = _resolve_safe_path(vault_root, path)
|
|
except ServiceError:
|
|
await websocket.close(code=4404)
|
|
return
|
|
|
|
if not file_path.exists() or not file_path.is_file():
|
|
await websocket.close(code=4404)
|
|
return
|
|
|
|
await websocket.accept()
|
|
await collab_manager.connect(websocket, vault_name, path, file_path, user)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Static files & SPA fallback
|
|
# ---------------------------------------------------------------------------
|
|
|
|
if FRONTEND_DIR.exists():
|
|
# ``Cache-Control`` for /static is set by SecurityHeadersMiddleware (no-cache).
|
|
app.mount("/static", StaticFiles(directory=str(FRONTEND_DIR)), name="static")
|
|
|
|
@app.get("/sw.js")
|
|
async def serve_service_worker():
|
|
"""Serve the service worker for PWA support."""
|
|
sw_file = FRONTEND_DIR / "sw.js"
|
|
if sw_file.exists():
|
|
return FileResponse(
|
|
sw_file,
|
|
media_type="application/javascript",
|
|
headers={
|
|
"Cache-Control": "no-cache, no-store, must-revalidate",
|
|
"Service-Worker-Allowed": "/"
|
|
}
|
|
)
|
|
raise HTTPException(status_code=404, detail="Service worker not found")
|
|
|
|
@app.get("/manifest.json")
|
|
async def serve_manifest():
|
|
"""Serve the PWA manifest."""
|
|
manifest_file = FRONTEND_DIR / "manifest.json"
|
|
if manifest_file.exists():
|
|
return FileResponse(
|
|
manifest_file,
|
|
media_type="application/manifest+json",
|
|
headers={"Cache-Control": "no-cache"}
|
|
)
|
|
raise HTTPException(status_code=404, detail="Manifest not found")
|
|
|
|
@app.get("/popout/{vault_name}/{path:path}")
|
|
async def serve_popout(vault_name: str, path: str):
|
|
"""Serve the minimalist popout page for a specific file."""
|
|
popout_file = FRONTEND_DIR / "popout.html"
|
|
if popout_file.exists():
|
|
return HTMLResponse(content=popout_file.read_text(encoding="utf-8"), headers={"Cache-Control": "no-cache"})
|
|
raise HTTPException(status_code=404, detail="Popout template not found")
|
|
|
|
@app.get("/editor-poc")
|
|
async def serve_editor_poc():
|
|
"""Serve the standalone Editor POC page (multi-zone toolbar demo)."""
|
|
poc_file = FRONTEND_DIR / "editor-poc.html"
|
|
if poc_file.exists():
|
|
return HTMLResponse(content=poc_file.read_text(encoding="utf-8"), headers={"Cache-Control": "no-cache"})
|
|
raise HTTPException(status_code=404, detail="Editor POC not found")
|
|
|
|
@app.get("/admin.html", response_class=HTMLResponse)
|
|
async def serve_admin_page(_current_user=Depends(require_admin)):
|
|
"""Serve the admin dashboard page (ROADMAP #71) — admin-gated.
|
|
|
|
Must be declared BEFORE the SPA catch-all ``/{full_path:path}`` or the
|
|
admin page would be shadowed by ``index.html`` (the reported bug: the
|
|
Admin menu kept returning to the main page).
|
|
"""
|
|
admin_file = FRONTEND_DIR / "admin.html"
|
|
if admin_file.exists():
|
|
return HTMLResponse(content=admin_file.read_text(encoding="utf-8"), headers={"Cache-Control": "no-cache"})
|
|
raise HTTPException(status_code=404, detail="Admin page not found")
|
|
|
|
@app.get("/{full_path:path}")
|
|
async def serve_spa(full_path: str):
|
|
"""Serve the SPA index.html for all non-API routes."""
|
|
index_file = FRONTEND_DIR / "index.html"
|
|
if index_file.exists():
|
|
return HTMLResponse(content=index_file.read_text(encoding="utf-8"), headers={"Cache-Control": "no-cache"})
|
|
raise HTTPException(status_code=404, detail="Frontend not found")
|