From 2c12af95fc92eda7f1cba051c399e39de068cf02 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Sat, 7 Mar 2026 10:00:47 +0000
Subject: [PATCH 1/7] Initial plan
From 89e0c2fb5055298aad40b9b493d37b95a7c0c480 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Sat, 7 Mar 2026 10:32:00 +0000
Subject: [PATCH 2/7] feat(pipelines): add custom processing pipeline engine
- Add Pipeline and PipelineStep models with user-specific ownership
- Add pipeline_id FK column to FileRecord
- Migration 017_add_pipelines (batch mode for SQLite FK compat)
- Pipeline CRUD API at /api/pipelines with step management endpoints
- Reorder steps PUT endpoint placed before parameterised {step_id} routes
- POST /api/files/{id}/assign-pipeline for per-file pipeline assignment
- Admin-only POST /api/pipelines/admin/system for system-level pipelines
- Management UI at /pipelines (Jinja2 + Alpine.js + Tailwind)
- Pipelines link added to desktop and mobile navigation
- 41 new tests in tests/test_api_pipelines.py (all passing)
Co-authored-by: christianlouis <361235+christianlouis@users.noreply.github.com>
---
app/api/__init__.py | 2 +
app/api/files.py | 72 ++
app/api/pipelines.py | 793 +++++++++++++++++++++
app/models.py | 70 ++
app/views/__init__.py | 2 +
app/views/pipelines.py | 34 +
frontend/templates/base.html | 2 +
frontend/templates/pipelines.html | 857 +++++++++++++++++++++++
migrations/versions/017_add_pipelines.py | 68 ++
tests/conftest.py | 2 +
tests/test_api_pipelines.py | 467 ++++++++++++
11 files changed, 2369 insertions(+)
create mode 100644 app/api/pipelines.py
create mode 100644 app/views/pipelines.py
create mode 100644 frontend/templates/pipelines.html
create mode 100644 migrations/versions/017_add_pipelines.py
create mode 100644 tests/test_api_pipelines.py
diff --git a/app/api/__init__.py b/app/api/__init__.py
index ac916de2..f266a266 100644
--- a/app/api/__init__.py
+++ b/app/api/__init__.py
@@ -17,6 +17,7 @@ from app.api.google_drive import router as google_drive_router
from app.api.logs import router as logs_router
from app.api.onedrive import router as onedrive_router
from app.api.openai import router as openai_router
+from app.api.pipelines import router as pipelines_router
from app.api.plans import router as plans_router
from app.api.process import router as process_router
from app.api.queue import router as queue_router
@@ -60,3 +61,4 @@ router.include_router(webhooks_router)
router.include_router(database_router)
router.include_router(subscriptions_router)
router.include_router(plans_router)
+router.include_router(pipelines_router)
diff --git a/app/api/files.py b/app/api/files.py
index 87586664..ee5d2998 100644
--- a/app/api/files.py
+++ b/app/api/files.py
@@ -1559,3 +1559,75 @@ def assign_owner(request: Request, db: DbSession, owner_id: str = Query(...), fi
"updated_count": updated,
"owner_id": owner_id,
}
+
+
+# ---------------------------------------------------------------------------
+# Pipeline assignment
+# ---------------------------------------------------------------------------
+
+
+@router.post("/files/{file_id}/assign-pipeline")
+@require_login
+def assign_pipeline_to_file(
+ request: Request,
+ file_id: int,
+ db: DbSession,
+ pipeline_id: int | None = None,
+):
+ """Assign (or remove) a processing pipeline from a file.
+
+ Path Parameters:
+ file_id: The file to update.
+
+ Query / Body Parameters:
+ pipeline_id: The pipeline to assign. Pass ``null`` or omit to clear the
+ assignment (the system default will be used for future processing).
+
+ Returns:
+ A summary dict with the file_id and updated pipeline_id.
+
+ Raises:
+ HTTPException 404: If the file or pipeline does not exist / is not
+ accessible to the current user.
+ """
+ from app.auth import get_current_user
+ from app.models import Pipeline
+
+ user = get_current_user(request)
+ # Derive user identity the same way the pipelines API does (_get_user_id)
+ if user:
+ user_id: str = user.get("preferred_username") or user.get("email") or user.get("id") or "anonymous"
+ else:
+ user_id = "anonymous"
+
+ is_admin_user = bool(user and user.get("is_admin"))
+
+ file_record = db.query(FileRecord).filter(FileRecord.id == file_id).first()
+ if not file_record:
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="File not found")
+
+ # Non-admins may only update files they own (or unowned files in single-user mode)
+ owner_id = get_current_owner_id(request)
+ if not is_admin_user and file_record.owner_id is not None and file_record.owner_id != owner_id:
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="File not found")
+
+ if pipeline_id is not None:
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline:
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+ # Check access: users can only assign their own pipelines or system pipelines (owner_id=None)
+ if not is_admin_user and pipeline.owner_id is not None and pipeline.owner_id != user_id:
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ file_record.pipeline_id = pipeline_id
+
+ try:
+ db.commit()
+ db.refresh(file_record)
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to assign pipeline to file id={file_id}: {exc}")
+ raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail="Failed to assign pipeline")
+
+ logger.info(f"Pipeline {pipeline_id!r} assigned to file id={file_id}")
+ return {"file_id": file_id, "pipeline_id": file_record.pipeline_id}
diff --git a/app/api/pipelines.py b/app/api/pipelines.py
new file mode 100644
index 00000000..13182a80
--- /dev/null
+++ b/app/api/pipelines.py
@@ -0,0 +1,793 @@
+"""
+Pipelines API endpoints.
+
+Provides full CRUD for processing pipelines and their steps. Pipelines are
+user-specific: regular users can only manage their own pipelines, while admins
+can also create and manage *system default* pipelines (owner_id = NULL) that
+are visible to all users.
+
+Built-in step types are exposed via GET /api/pipelines/step-types so that UIs
+can render the correct configuration form without hard-coding the catalogue.
+"""
+
+import json
+import logging
+from typing import Annotated, Any
+
+from fastapi import APIRouter, Body, Depends, HTTPException, Request, status
+from pydantic import BaseModel, Field
+from sqlalchemy.orm import Session
+
+from app.auth import get_current_user, require_login
+from app.database import get_db
+from app.models import Pipeline, PipelineStep
+
+logger = logging.getLogger(__name__)
+
+router = APIRouter(prefix="/pipelines", tags=["pipelines"])
+
+DbSession = Annotated[Session, Depends(get_db)]
+
+# ---------------------------------------------------------------------------
+# Built-in step type catalogue
+# ---------------------------------------------------------------------------
+
+PIPELINE_STEP_TYPES: dict[str, dict[str, Any]] = {
+ "convert_to_pdf": {
+ "label": "Convert to PDF",
+ "description": "Convert non-PDF documents to PDF format using Gotenberg.",
+ "config_schema": {},
+ },
+ "check_duplicates": {
+ "label": "Check for Duplicates",
+ "description": "Compare file hash against existing documents to detect duplicates.",
+ "config_schema": {},
+ },
+ "ocr": {
+ "label": "OCR Processing",
+ "description": "Extract text using Azure Document Intelligence or local Tesseract.",
+ "config_schema": {
+ "force_cloud_ocr": {
+ "type": "boolean",
+ "default": False,
+ "description": "Always use cloud OCR even if the PDF already has embedded text.",
+ }
+ },
+ },
+ "extract_metadata": {
+ "label": "Metadata Extraction",
+ "description": "Extract structured metadata (document type, sender, recipient, tags) using AI.",
+ "config_schema": {},
+ },
+ "embed_metadata": {
+ "label": "Embed Metadata into PDF",
+ "description": "Write the extracted metadata into the PDF document properties.",
+ "config_schema": {},
+ },
+ "compute_embedding": {
+ "label": "Compute Text Embedding",
+ "description": "Compute semantic text embeddings for full-text and similarity search.",
+ "config_schema": {},
+ },
+ "send_to_destinations": {
+ "label": "Send to Storage Destinations",
+ "description": "Upload the processed document to all configured storage destinations.",
+ "config_schema": {},
+ },
+ "classify": {
+ "label": "Document Classification",
+ "description": "Classify the document type using AI without full metadata extraction.",
+ "config_schema": {},
+ },
+}
+
+# ---------------------------------------------------------------------------
+# Helpers
+# ---------------------------------------------------------------------------
+
+MAX_STEPS_PER_PIPELINE = 50
+MAX_NAME_LENGTH = 255
+
+
+def _get_user_id(request: Request) -> str:
+ """Return a stable user identifier from the session."""
+ user = get_current_user(request)
+ if user:
+ return user.get("preferred_username") or user.get("email") or user.get("id", "anonymous")
+ return "anonymous"
+
+
+def _is_admin(request: Request) -> bool:
+ """Return True if the current session user is an admin."""
+ user = get_current_user(request)
+ return bool(user and user.get("is_admin"))
+
+
+def _can_access_pipeline(pipeline: Pipeline, user_id: str, admin: bool) -> bool:
+ """Return True if the user may read or write this pipeline."""
+ # System pipelines (owner_id=NULL) are readable by everyone; only admins can write
+ if pipeline.owner_id is None:
+ return True
+ # Own pipeline
+ return pipeline.owner_id == user_id or admin
+
+
+def _can_write_pipeline(pipeline: Pipeline, user_id: str, admin: bool) -> bool:
+ """Return True if the user may create/update/delete this pipeline."""
+ if pipeline.owner_id is None:
+ return admin
+ return pipeline.owner_id == user_id or admin
+
+
+def _serialize_step(step: PipelineStep) -> dict[str, Any]:
+ return {
+ "id": step.id,
+ "pipeline_id": step.pipeline_id,
+ "position": step.position,
+ "step_type": step.step_type,
+ "label": step.label,
+ "config": json.loads(step.config) if step.config else {},
+ "enabled": step.enabled,
+ "created_at": step.created_at.isoformat() if step.created_at else None,
+ "updated_at": step.updated_at.isoformat() if step.updated_at else None,
+ }
+
+
+def _serialize_pipeline(pipeline: Pipeline, include_steps: bool = False, db: Session | None = None) -> dict[str, Any]:
+ data: dict[str, Any] = {
+ "id": pipeline.id,
+ "owner_id": pipeline.owner_id,
+ "name": pipeline.name,
+ "description": pipeline.description,
+ "is_default": pipeline.is_default,
+ "is_active": pipeline.is_active,
+ "created_at": pipeline.created_at.isoformat() if pipeline.created_at else None,
+ "updated_at": pipeline.updated_at.isoformat() if pipeline.updated_at else None,
+ }
+ if include_steps and db is not None:
+ steps = (
+ db.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline.id).order_by(PipelineStep.position).all()
+ )
+ data["steps"] = [_serialize_step(s) for s in steps]
+ return data
+
+
+# ---------------------------------------------------------------------------
+# Pydantic schemas
+# ---------------------------------------------------------------------------
+
+
+class PipelineCreate(BaseModel):
+ """Body for creating a pipeline."""
+
+ name: str = Field(..., max_length=MAX_NAME_LENGTH, description="Human-readable pipeline name")
+ description: str | None = Field(default=None, max_length=4096)
+ is_default: bool = Field(default=False)
+ is_active: bool = Field(default=True)
+
+
+class PipelineUpdate(BaseModel):
+ """Body for updating a pipeline (all fields optional)."""
+
+ name: str | None = Field(default=None, max_length=MAX_NAME_LENGTH)
+ description: str | None = Field(default=None, max_length=4096)
+ is_default: bool | None = None
+ is_active: bool | None = None
+
+
+class PipelineStepCreate(BaseModel):
+ """Body for adding a step to a pipeline."""
+
+ step_type: str = Field(..., description="One of the recognised step type keys")
+ label: str | None = Field(default=None, max_length=MAX_NAME_LENGTH)
+ config: dict[str, Any] = Field(default_factory=dict)
+ enabled: bool = Field(default=True)
+ position: int | None = Field(default=None, ge=0, description="Insertion position; appended at end if omitted")
+
+
+class PipelineStepUpdate(BaseModel):
+ """Body for updating a pipeline step (all fields optional)."""
+
+ step_type: str | None = None
+ label: str | None = Field(default=None, max_length=MAX_NAME_LENGTH)
+ config: dict[str, Any] | None = None
+ enabled: bool | None = None
+ position: int | None = Field(default=None, ge=0)
+
+
+# ---------------------------------------------------------------------------
+# Step-types catalogue endpoint (no auth required — it's public metadata)
+# ---------------------------------------------------------------------------
+
+
+@router.get("/step-types")
+def list_step_types() -> dict[str, Any]:
+ """Return the catalogue of built-in pipeline step types.
+
+ Returns:
+ A mapping of step_type key → metadata (label, description, config_schema).
+ """
+ return PIPELINE_STEP_TYPES
+
+
+# ---------------------------------------------------------------------------
+# Pipeline CRUD
+# ---------------------------------------------------------------------------
+
+
+@router.get("")
+@require_login
+def list_pipelines(request: Request, db: DbSession) -> list[dict[str, Any]]:
+ """List pipelines visible to the current user.
+
+ Regular users see: their own pipelines + system pipelines (owner_id=NULL).
+ Admins see: all pipelines from all users.
+
+ Returns:
+ A list of pipeline objects (without steps — use GET /pipelines/{id} for steps).
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ if admin:
+ pipelines = db.query(Pipeline).order_by(Pipeline.owner_id.nullsfirst(), Pipeline.name).all()
+ else:
+ pipelines = (
+ db.query(Pipeline)
+ .filter((Pipeline.owner_id == user_id) | (Pipeline.owner_id.is_(None)))
+ .order_by(Pipeline.owner_id.nullsfirst(), Pipeline.name)
+ .all()
+ )
+
+ return [_serialize_pipeline(p) for p in pipelines]
+
+
+@router.post("", status_code=status.HTTP_201_CREATED)
+@require_login
+def create_pipeline(request: Request, db: DbSession, body: PipelineCreate) -> dict[str, Any]:
+ """Create a new pipeline for the current user.
+
+ Admins can create system default pipelines by passing ``owner_id=null``
+ via the body — however, that is handled implicitly: to create a system
+ pipeline, call ``POST /api/admin/pipelines`` (admin endpoint) instead.
+ Regular users always get their own user_id as owner.
+
+ Returns:
+ The created pipeline object.
+
+ Raises:
+ HTTPException 409: If a pipeline with the same name already exists for this owner.
+ """
+ user_id = _get_user_id(request)
+
+ name = body.name.strip() if body.name else ""
+ if not name:
+ raise HTTPException(
+ status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
+ detail="name is required",
+ )
+
+ # Enforce unique name per owner
+ existing = db.query(Pipeline).filter(Pipeline.owner_id == user_id, Pipeline.name == name).first()
+ if existing:
+ raise HTTPException(
+ status_code=status.HTTP_409_CONFLICT,
+ detail=f"A pipeline named '{name}' already exists",
+ )
+
+ # If this pipeline is marked as default, unset the existing default for this user
+ if body.is_default:
+ _unset_default(db, user_id)
+
+ pipeline = Pipeline(
+ owner_id=user_id,
+ name=name,
+ description=body.description,
+ is_default=body.is_default,
+ is_active=body.is_active,
+ )
+ try:
+ db.add(pipeline)
+ db.commit()
+ db.refresh(pipeline)
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to create pipeline user={user_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to create pipeline",
+ )
+
+ logger.info(f"Pipeline created: id={pipeline.id}, owner={user_id}, name={name!r}")
+ return _serialize_pipeline(pipeline)
+
+
+@router.get("/{pipeline_id}")
+@require_login
+def get_pipeline(pipeline_id: int, request: Request, db: DbSession) -> dict[str, Any]:
+ """Return a single pipeline with its steps.
+
+ Path Parameters:
+ pipeline_id: The ID of the pipeline.
+
+ Returns:
+ The pipeline object including its ordered steps.
+
+ Raises:
+ HTTPException 404: If the pipeline does not exist or is not accessible.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ return _serialize_pipeline(pipeline, include_steps=True, db=db)
+
+
+@router.put("/{pipeline_id}")
+@require_login
+def update_pipeline(pipeline_id: int, request: Request, db: DbSession, body: PipelineUpdate) -> dict[str, Any]:
+ """Update a pipeline's metadata.
+
+ Path Parameters:
+ pipeline_id: The ID of the pipeline to update.
+
+ Returns:
+ The updated pipeline object.
+
+ Raises:
+ HTTPException 403: If the caller does not own this pipeline.
+ HTTPException 404: If the pipeline does not exist.
+ HTTPException 409: If the new name conflicts with an existing pipeline.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ if not _can_write_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Cannot modify this pipeline")
+
+ if body.name is not None:
+ new_name = body.name.strip()
+ if not new_name:
+ raise HTTPException(
+ status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
+ detail="name must not be empty",
+ )
+ if new_name != pipeline.name:
+ conflict = (
+ db.query(Pipeline)
+ .filter(Pipeline.owner_id == pipeline.owner_id, Pipeline.name == new_name, Pipeline.id != pipeline_id)
+ .first()
+ )
+ if conflict:
+ raise HTTPException(
+ status_code=status.HTTP_409_CONFLICT,
+ detail=f"A pipeline named '{new_name}' already exists",
+ )
+ pipeline.name = new_name
+
+ if body.description is not None:
+ pipeline.description = body.description
+
+ if body.is_active is not None:
+ pipeline.is_active = body.is_active
+
+ if body.is_default is not None:
+ if body.is_default and not pipeline.is_default:
+ _unset_default(db, pipeline.owner_id)
+ pipeline.is_default = body.is_default
+
+ try:
+ db.commit()
+ db.refresh(pipeline)
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to update pipeline id={pipeline_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to update pipeline",
+ )
+
+ logger.info(f"Pipeline updated: id={pipeline_id}, user={user_id}")
+ return _serialize_pipeline(pipeline, include_steps=True, db=db)
+
+
+@router.delete("/{pipeline_id}", status_code=status.HTTP_204_NO_CONTENT)
+@require_login
+def delete_pipeline(pipeline_id: int, request: Request, db: DbSession) -> None:
+ """Delete a pipeline and all its steps.
+
+ Path Parameters:
+ pipeline_id: The ID of the pipeline to delete.
+
+ Raises:
+ HTTPException 403: If the caller does not own this pipeline.
+ HTTPException 404: If the pipeline does not exist.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ if not _can_write_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Cannot delete this pipeline")
+
+ try:
+ db.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline_id).delete()
+ db.delete(pipeline)
+ db.commit()
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to delete pipeline id={pipeline_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to delete pipeline",
+ )
+
+ logger.info(f"Pipeline deleted: id={pipeline_id}, user={user_id}")
+
+
+# ---------------------------------------------------------------------------
+# Admin-only: create system (owner_id=NULL) pipeline
+# ---------------------------------------------------------------------------
+
+
+@router.post("/admin/system", status_code=status.HTTP_201_CREATED, tags=["admin-pipelines"])
+@require_login
+def create_system_pipeline(request: Request, db: DbSession, body: PipelineCreate) -> dict[str, Any]:
+ """Create a system-level (owner_id=NULL) default pipeline. Admin only.
+
+ System pipelines are visible to all users and can be set as the global
+ default. Only admins may create them.
+
+ Returns:
+ The created system pipeline.
+
+ Raises:
+ HTTPException 403: If the caller is not an admin.
+ HTTPException 409: If a system pipeline with the same name already exists.
+ """
+ if not _is_admin(request):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Admin access required")
+
+ name = body.name.strip() if body.name else ""
+ if not name:
+ raise HTTPException(
+ status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
+ detail="name is required",
+ )
+
+ existing = db.query(Pipeline).filter(Pipeline.owner_id.is_(None), Pipeline.name == name).first()
+ if existing:
+ raise HTTPException(
+ status_code=status.HTTP_409_CONFLICT,
+ detail=f"A system pipeline named '{name}' already exists",
+ )
+
+ if body.is_default:
+ _unset_default(db, None)
+
+ pipeline = Pipeline(
+ owner_id=None,
+ name=name,
+ description=body.description,
+ is_default=body.is_default,
+ is_active=body.is_active,
+ )
+ try:
+ db.add(pipeline)
+ db.commit()
+ db.refresh(pipeline)
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to create system pipeline: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to create system pipeline",
+ )
+
+ logger.info(f"System pipeline created: id={pipeline.id}, name={name!r}")
+ return _serialize_pipeline(pipeline)
+
+
+# ---------------------------------------------------------------------------
+# Step management
+# ---------------------------------------------------------------------------
+
+
+@router.post("/{pipeline_id}/steps", status_code=status.HTTP_201_CREATED)
+@require_login
+def add_step(pipeline_id: int, request: Request, db: DbSession, body: PipelineStepCreate) -> dict[str, Any]:
+ """Add a step to a pipeline.
+
+ Steps are automatically appended at the end unless an explicit ``position``
+ is supplied. All existing steps at or after the insertion position are
+ shifted forward by one.
+
+ Path Parameters:
+ pipeline_id: The pipeline to add the step to.
+
+ Returns:
+ The created step object.
+
+ Raises:
+ HTTPException 403: If the caller cannot modify this pipeline.
+ HTTPException 404: If the pipeline does not exist.
+ HTTPException 422: If the step_type is not recognised.
+ HTTPException 409: If the maximum number of steps per pipeline is reached.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ if not _can_write_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Cannot modify this pipeline")
+
+ if body.step_type not in PIPELINE_STEP_TYPES:
+ raise HTTPException(
+ status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
+ detail=f"Unknown step type '{body.step_type}'. Valid types: {sorted(PIPELINE_STEP_TYPES)}",
+ )
+
+ current_count = db.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline_id).count()
+ if current_count >= MAX_STEPS_PER_PIPELINE:
+ raise HTTPException(
+ status_code=status.HTTP_409_CONFLICT,
+ detail=f"Maximum of {MAX_STEPS_PER_PIPELINE} steps per pipeline reached",
+ )
+
+ # Determine insertion position
+ if body.position is None:
+ max_pos = (
+ db.query(PipelineStep.position)
+ .filter(PipelineStep.pipeline_id == pipeline_id)
+ .order_by(PipelineStep.position.desc())
+ .first()
+ )
+ insert_pos = (max_pos[0] + 1) if max_pos else 0
+ else:
+ insert_pos = body.position
+ # Shift existing steps
+ steps_to_shift = (
+ db.query(PipelineStep)
+ .filter(PipelineStep.pipeline_id == pipeline_id, PipelineStep.position >= insert_pos)
+ .all()
+ )
+ for s in steps_to_shift:
+ s.position += 1
+
+ step = PipelineStep(
+ pipeline_id=pipeline_id,
+ position=insert_pos,
+ step_type=body.step_type,
+ label=body.label,
+ config=json.dumps(body.config) if body.config else None,
+ enabled=body.enabled,
+ )
+ try:
+ db.add(step)
+ db.commit()
+ db.refresh(step)
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to add step to pipeline id={pipeline_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to add step",
+ )
+
+ logger.info(f"Step added: pipeline={pipeline_id}, step_type={body.step_type!r}, pos={insert_pos}")
+ return _serialize_step(step)
+
+
+@router.put("/{pipeline_id}/steps/reorder", tags=["pipelines"])
+@require_login
+def reorder_steps(
+ pipeline_id: int,
+ request: Request,
+ db: DbSession,
+ step_ids: list[int] = Body(..., description="Ordered list of step IDs representing the new order"),
+) -> list[dict[str, Any]]:
+ """Replace the step order for a pipeline.
+
+ Provide a complete ordered list of *all* step IDs. Their ``position``
+ values will be reassigned 0, 1, 2, … in the given order.
+
+ Path Parameters:
+ pipeline_id: The pipeline whose steps are being reordered.
+
+ Returns:
+ The updated, ordered list of step objects.
+
+ Raises:
+ HTTPException 422: If the provided list does not contain exactly the
+ current set of step IDs for this pipeline.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ if not _can_write_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Cannot modify this pipeline")
+
+ existing_steps = db.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline_id).all()
+ existing_ids = {s.id for s in existing_steps}
+
+ if set(step_ids) != existing_ids or len(step_ids) != len(existing_ids):
+ raise HTTPException(
+ status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
+ detail="step_ids must contain exactly the current step IDs for this pipeline",
+ )
+
+ step_map = {s.id: s for s in existing_steps}
+ for pos, sid in enumerate(step_ids):
+ step_map[sid].position = pos
+
+ try:
+ db.commit()
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to reorder steps for pipeline id={pipeline_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to reorder steps",
+ )
+
+ updated = (
+ db.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline_id).order_by(PipelineStep.position).all()
+ )
+ return [_serialize_step(s) for s in updated]
+
+
+@router.put("/{pipeline_id}/steps/{step_id}")
+@require_login
+def update_step(
+ pipeline_id: int, step_id: int, request: Request, db: DbSession, body: PipelineStepUpdate
+) -> dict[str, Any]:
+ """Update an existing pipeline step.
+
+ Path Parameters:
+ pipeline_id: The owning pipeline.
+ step_id: The step to update.
+
+ Returns:
+ The updated step object.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ if not _can_write_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Cannot modify this pipeline")
+
+ step = db.query(PipelineStep).filter(PipelineStep.id == step_id, PipelineStep.pipeline_id == pipeline_id).first()
+ if not step:
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Step not found")
+
+ if body.step_type is not None:
+ if body.step_type not in PIPELINE_STEP_TYPES:
+ raise HTTPException(
+ status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
+ detail=f"Unknown step type '{body.step_type}'",
+ )
+ step.step_type = body.step_type
+
+ if body.label is not None:
+ step.label = body.label
+
+ if body.config is not None:
+ step.config = json.dumps(body.config)
+
+ if body.enabled is not None:
+ step.enabled = body.enabled
+
+ if body.position is not None and body.position != step.position:
+ old_pos = step.position
+ new_pos = body.position
+ if new_pos > old_pos:
+ # Moving down: shift intervening steps up
+ db.query(PipelineStep).filter(
+ PipelineStep.pipeline_id == pipeline_id,
+ PipelineStep.position > old_pos,
+ PipelineStep.position <= new_pos,
+ PipelineStep.id != step_id,
+ ).update({"position": PipelineStep.position - 1})
+ else:
+ # Moving up: shift intervening steps down
+ db.query(PipelineStep).filter(
+ PipelineStep.pipeline_id == pipeline_id,
+ PipelineStep.position >= new_pos,
+ PipelineStep.position < old_pos,
+ PipelineStep.id != step_id,
+ ).update({"position": PipelineStep.position + 1})
+ step.position = new_pos
+
+ try:
+ db.commit()
+ db.refresh(step)
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to update step id={step_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to update step",
+ )
+
+ return _serialize_step(step)
+
+
+@router.delete("/{pipeline_id}/steps/{step_id}", status_code=status.HTTP_204_NO_CONTENT)
+@require_login
+def delete_step(pipeline_id: int, step_id: int, request: Request, db: DbSession) -> None:
+ """Delete a step from a pipeline.
+
+ Path Parameters:
+ pipeline_id: The owning pipeline.
+ step_id: The step to delete.
+ """
+ user_id = _get_user_id(request)
+ admin = _is_admin(request)
+
+ pipeline = db.query(Pipeline).filter(Pipeline.id == pipeline_id).first()
+ if not pipeline or not _can_access_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Pipeline not found")
+
+ if not _can_write_pipeline(pipeline, user_id, admin):
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Cannot modify this pipeline")
+
+ step = db.query(PipelineStep).filter(PipelineStep.id == step_id, PipelineStep.pipeline_id == pipeline_id).first()
+ if not step:
+ raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Step not found")
+
+ deleted_pos = step.position
+ try:
+ db.delete(step)
+ # Compact remaining step positions
+ db.query(PipelineStep).filter(
+ PipelineStep.pipeline_id == pipeline_id,
+ PipelineStep.position > deleted_pos,
+ ).update({"position": PipelineStep.position - 1})
+ db.commit()
+ except Exception as exc:
+ db.rollback()
+ logger.exception(f"Failed to delete step id={step_id}: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to delete step",
+ )
+
+ logger.info(f"Step deleted: id={step_id}, pipeline={pipeline_id}")
+
+
+# ---------------------------------------------------------------------------
+# Helper: unset default flag for an owner
+# ---------------------------------------------------------------------------
+
+
+def _unset_default(db: Session, owner_id: str | None) -> None:
+ """Clear the is_default flag on all pipelines for the given owner."""
+ if owner_id is None:
+ db.query(Pipeline).filter(Pipeline.owner_id.is_(None), Pipeline.is_default.is_(True)).update(
+ {"is_default": False}
+ )
+ else:
+ db.query(Pipeline).filter(Pipeline.owner_id == owner_id, Pipeline.is_default.is_(True)).update(
+ {"is_default": False}
+ )
diff --git a/app/models.py b/app/models.py
index b5857bb5..1287ff40 100644
--- a/app/models.py
+++ b/app/models.py
@@ -6,6 +6,7 @@ from app.database import Base
# Foreign key constants
_FILES_ID_FK = "files.id"
+_PIPELINES_ID_FK = "pipelines.id"
class DocumentMetadata(Base):
@@ -79,6 +80,9 @@ class FileRecord(Base):
# Pre-computed text embedding vector stored as JSON array of floats
embedding = Column(Text, nullable=True)
+ # Processing pipeline assigned to this file (NULL = use system default)
+ pipeline_id = Column(Integer, ForeignKey(_PIPELINES_ID_FK), nullable=True, index=True)
+
# Timestamp when we inserted this record
created_at = Column(DateTime(timezone=True), server_default=func.now(), index=True)
@@ -259,3 +263,69 @@ class SubscriptionPlan(Base):
created_at = Column(DateTime(timezone=True), server_default=func.now())
updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now())
+
+
+class Pipeline(Base):
+ """User-defined processing pipeline: an ordered set of steps.
+
+ Pipelines are user-specific. A pipeline with ``owner_id = NULL`` is a
+ *system default* pipeline that only admins may create. Regular users
+ create pipelines under their own ``owner_id``. When a file has no
+ explicit pipeline assigned, the active system default is used.
+ """
+
+ __tablename__ = "pipelines"
+
+ id = Column(Integer, primary_key=True, index=True)
+
+ # Owner of this pipeline. NULL = system/admin pipeline visible to everyone.
+ owner_id = Column(String, nullable=True, index=True)
+
+ # Human-readable name (unique per owner)
+ name = Column(String(255), nullable=False)
+
+ # Optional description
+ description = Column(Text, nullable=True)
+
+ # When True this pipeline is the default for new files belonging to the owner
+ # (or the global default when owner_id is NULL). Only one pipeline per
+ # owner may be active default at a time — enforced at the application level.
+ is_default = Column(Boolean, nullable=False, default=False)
+
+ # Soft-disable without deleting
+ is_active = Column(Boolean, nullable=False, default=True)
+
+ created_at = Column(DateTime(timezone=True), server_default=func.now())
+ updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now())
+
+
+class PipelineStep(Base):
+ """A single step in a processing pipeline.
+
+ Steps are executed in ascending ``position`` order. Each step has a
+ ``step_type`` that maps to a built-in processing action and an optional
+ ``config`` JSON blob with step-specific parameters.
+ """
+
+ __tablename__ = "pipeline_steps"
+
+ id = Column(Integer, primary_key=True, index=True)
+ pipeline_id = Column(Integer, ForeignKey(_PIPELINES_ID_FK), nullable=False, index=True)
+
+ # Execution order within the pipeline (lower = earlier)
+ position = Column(Integer, nullable=False, default=0)
+
+ # One of the recognised step types (see PIPELINE_STEP_TYPES in pipelines.py)
+ step_type = Column(String(100), nullable=False)
+
+ # Optional human-readable label override (defaults to step_type label)
+ label = Column(String(255), nullable=True)
+
+ # JSON-encoded step-specific configuration dict
+ config = Column(Text, nullable=True)
+
+ # When False this step is skipped during execution
+ enabled = Column(Boolean, nullable=False, default=True)
+
+ created_at = Column(DateTime(timezone=True), server_default=func.now())
+ updated_at = Column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now())
diff --git a/app/views/__init__.py b/app/views/__init__.py
index 75eafbfb..309dd3d1 100644
--- a/app/views/__init__.py
+++ b/app/views/__init__.py
@@ -14,6 +14,7 @@ from app.views.general import router as general_router
from app.views.google_drive import router as google_drive_router
from app.views.license_routes import router as license_router # Add the license router
from app.views.onedrive import router as onedrive_router
+from app.views.pipelines import router as pipelines_router # Processing pipelines
from app.views.plans import router as plans_router # Admin Plan Designer
from app.views.queue import router as queue_router
from app.views.search import router as search_router
@@ -39,3 +40,4 @@ router.include_router(search_router)
router.include_router(queue_router)
router.include_router(subscriptions_router) # Pricing + subscription pages
router.include_router(plans_router) # Admin Plan Designer
+router.include_router(pipelines_router) # Processing pipelines
diff --git a/app/views/pipelines.py b/app/views/pipelines.py
new file mode 100644
index 00000000..a0619f7c
--- /dev/null
+++ b/app/views/pipelines.py
@@ -0,0 +1,34 @@
+"""Pipelines view: management UI for processing pipelines."""
+
+import logging
+
+from fastapi import HTTPException, Request, status
+
+from app.views.base import APIRouter, require_login, settings, templates
+
+logger = logging.getLogger(__name__)
+router = APIRouter()
+
+
+@router.get("/pipelines")
+@require_login
+async def pipelines_page(request: Request):
+ """Processing pipeline management page for the current user.
+
+ Regular users manage their own pipelines. Admins additionally have access
+ to system-level pipelines through the same UI.
+ """
+ try:
+ return templates.TemplateResponse(
+ "pipelines.html",
+ {
+ "request": request,
+ "app_version": settings.version,
+ },
+ )
+ except Exception as exc:
+ logger.error(f"Error loading pipelines page: {exc}")
+ raise HTTPException(
+ status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
+ detail="Failed to load pipelines page",
+ )
diff --git a/frontend/templates/base.html b/frontend/templates/base.html
index a6db72cb..49af2bad 100644
--- a/frontend/templates/base.html
+++ b/frontend/templates/base.html
@@ -63,6 +63,7 @@
UploadFilesSearch
+ PipelinesPricing
@@ -175,6 +176,7 @@
UploadFilesSearch
+ PipelinesPricing
diff --git a/frontend/templates/pipelines.html b/frontend/templates/pipelines.html
new file mode 100644
index 00000000..fb36ec47
--- /dev/null
+++ b/frontend/templates/pipelines.html
@@ -0,0 +1,857 @@
+{% extends "base.html" %}
+{% block title %}Processing Pipelines – DocuElevate{% endblock %}
+
+{% block content %}
+
+
+
+
+
+
+
+ Processing Pipelines
+
+
+ Define and manage custom document processing workflows.
+ System pipelines (created by admins) are shown with a
+ System
+ badge and are visible to all users.
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
Loading pipelines…
+
+
+
+
+
+
+
+
No pipelines yet
+
Create your first pipeline to define custom document processing workflows.
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+ System
+
+
+ Default
+
+
+ Inactive
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
No steps defined. Add a step to start building your pipeline.
+
+
+
+
+
+
+
+
+
+
+
+
+
+ Disabled
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+{% endblock %}
+
+{% block head_extra %}
+
+{% endblock %}
diff --git a/migrations/versions/017_add_pipelines.py b/migrations/versions/017_add_pipelines.py
new file mode 100644
index 00000000..d14abbf9
--- /dev/null
+++ b/migrations/versions/017_add_pipelines.py
@@ -0,0 +1,68 @@
+"""Add pipelines and pipeline_steps tables; add pipeline_id to files
+
+Revision ID: 017_add_pipelines
+Revises: 016_add_userprofile_billing
+Create Date: 2026-03-07
+
+"""
+
+from typing import Union
+
+import sqlalchemy as sa
+from alembic import op
+
+# revision identifiers, used by Alembic.
+revision: str = "017_add_pipelines"
+down_revision: Union[str, None] = "016_add_userprofile_billing"
+depends_on: Union[str, None] = None
+
+
+def upgrade() -> None:
+ """Create pipelines, pipeline_steps tables and add pipeline_id FK to files."""
+ op.create_table(
+ "pipelines",
+ sa.Column("id", sa.Integer(), primary_key=True, index=True),
+ sa.Column("owner_id", sa.String(), nullable=True, index=True),
+ sa.Column("name", sa.String(255), nullable=False),
+ sa.Column("description", sa.Text(), nullable=True),
+ sa.Column("is_default", sa.Boolean(), nullable=False, server_default="0"),
+ sa.Column("is_active", sa.Boolean(), nullable=False, server_default="1"),
+ sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now()),
+ sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now()),
+ )
+
+ op.create_table(
+ "pipeline_steps",
+ sa.Column("id", sa.Integer(), primary_key=True, index=True),
+ sa.Column("pipeline_id", sa.Integer(), sa.ForeignKey("pipelines.id"), nullable=False, index=True),
+ sa.Column("position", sa.Integer(), nullable=False, server_default="0"),
+ sa.Column("step_type", sa.String(100), nullable=False),
+ sa.Column("label", sa.String(255), nullable=True),
+ sa.Column("config", sa.Text(), nullable=True),
+ sa.Column("enabled", sa.Boolean(), nullable=False, server_default="1"),
+ sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now()),
+ sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now()),
+ )
+
+ # Use batch mode for SQLite compatibility when altering the files table.
+ # SQLite does not support adding FK constraints inline via ALTER TABLE, so we
+ # add the plain integer column and define the FK reference at the model level.
+ with op.batch_alter_table("files") as batch_op:
+ batch_op.add_column(
+ sa.Column("pipeline_id", sa.Integer(), nullable=True),
+ )
+ batch_op.create_foreign_key(
+ "fk_files_pipeline_id",
+ "pipelines",
+ ["pipeline_id"],
+ ["id"],
+ )
+
+
+def downgrade() -> None:
+ """Drop pipeline_id from files and remove pipeline tables."""
+ with op.batch_alter_table("files") as batch_op:
+ batch_op.drop_column("pipeline_id")
+
+ op.drop_table("pipeline_steps")
+ op.drop_table("pipelines")
diff --git a/tests/conftest.py b/tests/conftest.py
index 98cb5316..5f68bd9d 100644
--- a/tests/conftest.py
+++ b/tests/conftest.py
@@ -62,6 +62,8 @@ from app.main import app as fastapi_app # noqa: E402
from app.models import ( # noqa: F401, E402
DocumentMetadata,
FileRecord,
+ Pipeline,
+ PipelineStep,
ProcessingLog,
SavedSearch,
UserProfile,
diff --git a/tests/test_api_pipelines.py b/tests/test_api_pipelines.py
new file mode 100644
index 00000000..ad7fcf33
--- /dev/null
+++ b/tests/test_api_pipelines.py
@@ -0,0 +1,467 @@
+"""Tests for the pipelines API endpoints.
+
+Covers CRUD operations for pipelines and steps, ownership/admin access control,
+step reordering, and the assign-pipeline-to-file endpoint.
+"""
+
+import json
+from unittest.mock import patch
+
+import pytest
+
+from app.models import FileRecord, Pipeline, PipelineStep
+
+# ---------------------------------------------------------------------------
+# Helpers
+# ---------------------------------------------------------------------------
+
+
+def _make_file(db_session, owner_id=None):
+ """Insert a minimal FileRecord and return it.
+
+ Default owner_id=None so tests work without an authenticated session.
+ """
+ fr = FileRecord(
+ owner_id=owner_id,
+ filehash="abc123",
+ original_filename="test.pdf",
+ local_filename="/tmp/test.pdf",
+ file_size=1024,
+ mime_type="application/pdf",
+ )
+ db_session.add(fr)
+ db_session.commit()
+ db_session.refresh(fr)
+ return fr
+
+
+# ---------------------------------------------------------------------------
+# Unit tests – step-types catalogue
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.unit
+class TestStepTypesCatalogue:
+ """Tests for the step-types read-only catalogue endpoint."""
+
+ def test_step_types_returns_dict(self, client):
+ """GET /api/pipelines/step-types returns a dict of known types."""
+ r = client.get("/api/pipelines/step-types")
+ assert r.status_code == 200
+ data = r.json()
+ assert isinstance(data, dict)
+ # Core built-in types must be present
+ for key in ("convert_to_pdf", "ocr", "extract_metadata", "embed_metadata", "compute_embedding"):
+ assert key in data, f"Expected step type '{key}' in catalogue"
+
+ def test_each_type_has_label_and_description(self, client):
+ """Each step-type entry has at least a label and description."""
+ r = client.get("/api/pipelines/step-types")
+ for key, meta in r.json().items():
+ assert "label" in meta, f"Step type '{key}' missing 'label'"
+ assert "description" in meta, f"Step type '{key}' missing 'description'"
+
+
+# ---------------------------------------------------------------------------
+# Integration tests – Pipeline CRUD
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.integration
+class TestPipelineCRUD:
+ """Full CRUD test-suite for pipeline management."""
+
+ def test_list_pipelines_empty(self, client):
+ """List returns an empty array when no pipelines exist."""
+ r = client.get("/api/pipelines")
+ assert r.status_code == 200
+ assert r.json() == []
+
+ def test_create_pipeline(self, client):
+ """POST /api/pipelines creates a new pipeline."""
+ r = client.post(
+ "/api/pipelines",
+ json={"name": "My Pipeline", "description": "Test description"},
+ )
+ assert r.status_code == 201
+ data = r.json()
+ assert data["name"] == "My Pipeline"
+ assert data["description"] == "Test description"
+ assert data["is_default"] is False
+ assert data["is_active"] is True
+ assert data["id"] is not None
+
+ def test_create_pipeline_duplicate_name_same_owner_rejected(self, client):
+ """Creating two pipelines with the same name is rejected with 409."""
+ client.post("/api/pipelines", json={"name": "Dupe"})
+ r = client.post("/api/pipelines", json={"name": "Dupe"})
+ assert r.status_code == 409
+
+ def test_create_pipeline_empty_name_rejected(self, client):
+ """An empty pipeline name returns 422."""
+ r = client.post("/api/pipelines", json={"name": " "})
+ assert r.status_code == 422
+
+ def test_get_pipeline_includes_steps(self, client):
+ """GET /api/pipelines/{id} returns the pipeline with a steps array."""
+ created = client.post("/api/pipelines", json={"name": "With Steps"}).json()
+ r = client.get(f"/api/pipelines/{created['id']}")
+ assert r.status_code == 200
+ assert "steps" in r.json()
+ assert r.json()["steps"] == []
+
+ def test_get_pipeline_not_found(self, client):
+ """GET on a non-existent pipeline returns 404."""
+ r = client.get("/api/pipelines/99999")
+ assert r.status_code == 404
+
+ def test_update_pipeline(self, client):
+ """PUT /api/pipelines/{id} updates name and description."""
+ created = client.post("/api/pipelines", json={"name": "Original"}).json()
+ r = client.put(
+ f"/api/pipelines/{created['id']}",
+ json={"name": "Renamed", "description": "New desc"},
+ )
+ assert r.status_code == 200
+ assert r.json()["name"] == "Renamed"
+ assert r.json()["description"] == "New desc"
+
+ def test_update_pipeline_empty_name_rejected(self, client):
+ """Updating a pipeline with an empty name returns 422."""
+ created = client.post("/api/pipelines", json={"name": "Good"}).json()
+ r = client.put(f"/api/pipelines/{created['id']}", json={"name": ""})
+ assert r.status_code == 422
+
+ def test_update_pipeline_name_conflict_rejected(self, client):
+ """Updating a pipeline's name to one already taken returns 409."""
+ client.post("/api/pipelines", json={"name": "Taken"})
+ second = client.post("/api/pipelines", json={"name": "Other"}).json()
+ r = client.put(f"/api/pipelines/{second['id']}", json={"name": "Taken"})
+ assert r.status_code == 409
+
+ def test_delete_pipeline(self, client):
+ """DELETE /api/pipelines/{id} removes the pipeline."""
+ created = client.post("/api/pipelines", json={"name": "Deletable"}).json()
+ r = client.delete(f"/api/pipelines/{created['id']}")
+ assert r.status_code == 204
+ assert client.get(f"/api/pipelines/{created['id']}").status_code == 404
+
+ def test_delete_pipeline_not_found(self, client):
+ """Deleting a non-existent pipeline returns 404."""
+ r = client.delete("/api/pipelines/99999")
+ assert r.status_code == 404
+
+ def test_is_default_flag(self, client):
+ """Setting is_default=True marks the pipeline as default."""
+ r = client.post("/api/pipelines", json={"name": "Default Pipeline", "is_default": True})
+ assert r.status_code == 201
+ assert r.json()["is_default"] is True
+
+ def test_only_one_default_per_owner(self, client):
+ """When a new default is set, the old one is cleared."""
+ first = client.post("/api/pipelines", json={"name": "First Default", "is_default": True}).json()
+ second = client.post("/api/pipelines", json={"name": "Second Default", "is_default": True}).json()
+
+ assert second["is_default"] is True
+ # The first should no longer be default
+ first_updated = client.get(f"/api/pipelines/{first['id']}").json()
+ assert first_updated["is_default"] is False
+
+ def test_list_returns_created_pipeline(self, client):
+ """After creating a pipeline it appears in the list."""
+ client.post("/api/pipelines", json={"name": "Visible"})
+ r = client.get("/api/pipelines")
+ names = [p["name"] for p in r.json()]
+ assert "Visible" in names
+
+
+# ---------------------------------------------------------------------------
+# Integration tests – Step management
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.integration
+class TestPipelineStepManagement:
+ """Tests for adding, updating, deleting, and reordering pipeline steps."""
+
+ def _create_pipeline(self, client, name="Test Pipeline"):
+ return client.post("/api/pipelines", json={"name": name}).json()
+
+ def test_add_step(self, client):
+ """POST /api/pipelines/{id}/steps adds a step."""
+ p = self._create_pipeline(client)
+ r = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"})
+ assert r.status_code == 201
+ step = r.json()
+ assert step["step_type"] == "ocr"
+ assert step["position"] == 0
+ assert step["enabled"] is True
+
+ def test_add_step_with_config(self, client):
+ """A step can be added with a custom config dict."""
+ p = self._create_pipeline(client)
+ r = client.post(
+ f"/api/pipelines/{p['id']}/steps",
+ json={"step_type": "ocr", "config": {"force_cloud_ocr": True}},
+ )
+ assert r.status_code == 201
+ assert r.json()["config"]["force_cloud_ocr"] is True
+
+ def test_add_step_with_custom_label(self, client):
+ """A step can override the default label."""
+ p = self._create_pipeline(client)
+ r = client.post(
+ f"/api/pipelines/{p['id']}/steps",
+ json={"step_type": "convert_to_pdf", "label": "My Converter"},
+ )
+ assert r.status_code == 201
+ assert r.json()["label"] == "My Converter"
+
+ def test_add_invalid_step_type_rejected(self, client):
+ """An unrecognised step type returns 422."""
+ p = self._create_pipeline(client)
+ r = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "nonexistent_step"})
+ assert r.status_code == 422
+
+ def test_steps_appended_in_order(self, client):
+ """Multiple steps are appended in position order."""
+ p = self._create_pipeline(client)
+ client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "convert_to_pdf"})
+ client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"})
+ client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "extract_metadata"})
+
+ details = client.get(f"/api/pipelines/{p['id']}").json()
+ types = [s["step_type"] for s in details["steps"]]
+ assert types == ["convert_to_pdf", "ocr", "extract_metadata"]
+
+ def test_update_step(self, client):
+ """PUT /api/pipelines/{id}/steps/{step_id} updates enabled flag."""
+ p = self._create_pipeline(client)
+ step = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"}).json()
+
+ r = client.put(
+ f"/api/pipelines/{p['id']}/steps/{step['id']}",
+ json={"enabled": False},
+ )
+ assert r.status_code == 200
+ assert r.json()["enabled"] is False
+
+ def test_update_step_not_found(self, client):
+ """Updating a step on the wrong pipeline returns 404."""
+ p = self._create_pipeline(client)
+ r = client.put(f"/api/pipelines/{p['id']}/steps/99999", json={"enabled": False})
+ assert r.status_code == 404
+
+ def test_delete_step(self, client):
+ """DELETE /api/pipelines/{id}/steps/{step_id} removes the step."""
+ p = self._create_pipeline(client)
+ step = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"}).json()
+
+ r = client.delete(f"/api/pipelines/{p['id']}/steps/{step['id']}")
+ assert r.status_code == 204
+
+ details = client.get(f"/api/pipelines/{p['id']}").json()
+ assert details["steps"] == []
+
+ def test_delete_step_compacts_positions(self, client):
+ """After deleting a step, remaining steps have contiguous positions."""
+ p = self._create_pipeline(client)
+ s1 = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "convert_to_pdf"}).json()
+ client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"})
+ client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "extract_metadata"})
+
+ client.delete(f"/api/pipelines/{p['id']}/steps/{s1['id']}")
+
+ details = client.get(f"/api/pipelines/{p['id']}").json()
+ positions = [s["position"] for s in details["steps"]]
+ assert positions == sorted(positions)
+ assert positions[0] == 0
+
+ def test_reorder_steps(self, client):
+ """PUT /api/pipelines/{id}/steps/reorder reorders all steps."""
+ p = self._create_pipeline(client)
+ s1 = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "convert_to_pdf"}).json()
+ s2 = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"}).json()
+ s3 = client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "extract_metadata"}).json()
+
+ # Reverse order
+ r = client.put(
+ f"/api/pipelines/{p['id']}/steps/reorder",
+ json=[s3["id"], s2["id"], s1["id"]],
+ )
+ assert r.status_code == 200
+ types = [s["step_type"] for s in r.json()]
+ assert types == ["extract_metadata", "ocr", "convert_to_pdf"]
+
+ def test_reorder_steps_invalid_ids_rejected(self, client):
+ """Providing wrong step IDs returns 422."""
+ p = self._create_pipeline(client)
+ client.post(f"/api/pipelines/{p['id']}/steps", json={"step_type": "ocr"})
+
+ r = client.put(f"/api/pipelines/{p['id']}/steps/reorder", json=[99999])
+ assert r.status_code == 422
+
+
+# ---------------------------------------------------------------------------
+# Integration tests – System pipeline (admin endpoint)
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.integration
+class TestSystemPipeline:
+ """Tests for the admin system-pipeline creation endpoint."""
+
+ def test_create_system_pipeline_as_admin(self, client):
+ """An admin can create a system pipeline (owner_id=NULL)."""
+ with patch("app.api.pipelines._is_admin", return_value=True):
+ r = client.post(
+ "/api/pipelines/admin/system",
+ json={"name": "Global Default", "is_default": True},
+ )
+ assert r.status_code == 201
+ data = r.json()
+ assert data["owner_id"] is None
+ assert data["is_default"] is True
+
+ def test_create_system_pipeline_as_non_admin_forbidden(self, client):
+ """A non-admin user cannot create a system pipeline."""
+ r = client.post(
+ "/api/pipelines/admin/system",
+ json={"name": "Should Fail"},
+ )
+ assert r.status_code == 403
+
+
+# ---------------------------------------------------------------------------
+# Integration tests – File pipeline assignment
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.integration
+class TestAssignPipelineToFile:
+ """Tests for POST /api/files/{id}/assign-pipeline."""
+
+ def test_assign_pipeline_to_file(self, client, db_session):
+ """Assigning a pipeline to a file stores pipeline_id on the record."""
+ # File with no owner so anonymous test session can access it
+ fr = _make_file(db_session, owner_id=None)
+ pipeline = client.post("/api/pipelines", json={"name": "Assign Test"}).json()
+
+ r = client.post(f"/api/files/{fr.id}/assign-pipeline?pipeline_id={pipeline['id']}")
+ assert r.status_code == 200
+ data = r.json()
+ assert data["file_id"] == fr.id
+ assert data["pipeline_id"] == pipeline["id"]
+
+ db_session.refresh(fr)
+ assert fr.pipeline_id == pipeline["id"]
+
+ def test_clear_pipeline_from_file(self, client, db_session):
+ """Passing no pipeline_id clears the assignment."""
+ fr = _make_file(db_session, owner_id=None)
+ pipeline = client.post("/api/pipelines", json={"name": "Clearable"}).json()
+ client.post(f"/api/files/{fr.id}/assign-pipeline?pipeline_id={pipeline['id']}")
+
+ r = client.post(f"/api/files/{fr.id}/assign-pipeline")
+ assert r.status_code == 200
+ assert r.json()["pipeline_id"] is None
+
+ def test_assign_nonexistent_pipeline_returns_404(self, client, db_session):
+ """Assigning a non-existent pipeline returns 404."""
+ fr = _make_file(db_session, owner_id=None)
+ r = client.post(f"/api/files/{fr.id}/assign-pipeline?pipeline_id=99999")
+ assert r.status_code == 404
+
+ def test_assign_pipeline_to_nonexistent_file_returns_404(self, client):
+ """Assigning a pipeline to a non-existent file returns 404."""
+ r = client.post("/api/files/99999/assign-pipeline?pipeline_id=1")
+ assert r.status_code == 404
+
+
+# ---------------------------------------------------------------------------
+# Unit tests – API helper logic
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.unit
+class TestPipelineAPIHelpers:
+ """Unit tests for API helper functions (no DB required)."""
+
+ def test_serialize_step_returns_expected_keys(self, db_session):
+ """_serialize_step returns all required fields."""
+ from app.api.pipelines import _serialize_step
+
+ # Build a minimal in-memory step
+ p = Pipeline(owner_id="u1", name="P", is_default=False, is_active=True)
+ db_session.add(p)
+ db_session.commit()
+
+ s = PipelineStep(
+ pipeline_id=p.id,
+ position=0,
+ step_type="ocr",
+ label="OCR",
+ config=json.dumps({"force_cloud_ocr": False}),
+ enabled=True,
+ )
+ db_session.add(s)
+ db_session.commit()
+
+ result = _serialize_step(s)
+ for key in ("id", "pipeline_id", "position", "step_type", "label", "config", "enabled"):
+ assert key in result, f"Missing key '{key}' in serialized step"
+ assert result["config"] == {"force_cloud_ocr": False}
+
+ def test_serialize_pipeline_returns_expected_keys(self, db_session):
+ """_serialize_pipeline returns all required fields."""
+ from app.api.pipelines import _serialize_pipeline
+
+ p = Pipeline(owner_id="u1", name="MyPipeline", is_default=True, is_active=True)
+ db_session.add(p)
+ db_session.commit()
+
+ result = _serialize_pipeline(p)
+ for key in ("id", "owner_id", "name", "description", "is_default", "is_active"):
+ assert key in result, f"Missing key '{key}' in serialized pipeline"
+
+ def test_can_access_system_pipeline(self):
+ """Anyone can read a system pipeline (owner_id=None)."""
+ from app.api.pipelines import _can_access_pipeline
+
+ p = Pipeline(owner_id=None, name="System", is_default=False, is_active=True)
+ assert _can_access_pipeline(p, "any_user", admin=False) is True
+
+ def test_cannot_write_system_pipeline_as_regular_user(self):
+ """Regular users cannot modify system pipelines."""
+ from app.api.pipelines import _can_write_pipeline
+
+ p = Pipeline(owner_id=None, name="System", is_default=False, is_active=True)
+ assert _can_write_pipeline(p, "regular_user", admin=False) is False
+
+ def test_admin_can_write_system_pipeline(self):
+ """Admins can modify system pipelines."""
+ from app.api.pipelines import _can_write_pipeline
+
+ p = Pipeline(owner_id=None, name="System", is_default=False, is_active=True)
+ assert _can_write_pipeline(p, "admin", admin=True) is True
+
+ def test_user_can_access_own_pipeline(self):
+ """A user can access pipelines they own."""
+ from app.api.pipelines import _can_access_pipeline
+
+ p = Pipeline(owner_id="user1", name="Mine", is_default=False, is_active=True)
+ assert _can_access_pipeline(p, "user1", admin=False) is True
+
+ def test_user_cannot_access_other_users_pipeline(self):
+ """A regular user cannot access another user's pipeline."""
+ from app.api.pipelines import _can_access_pipeline
+
+ p = Pipeline(owner_id="user1", name="Theirs", is_default=False, is_active=True)
+ assert _can_access_pipeline(p, "user2", admin=False) is False
+
+ def test_admin_can_access_any_pipeline(self):
+ """Admins can access any pipeline regardless of owner."""
+ from app.api.pipelines import _can_access_pipeline
+
+ p = Pipeline(owner_id="user99", name="Private", is_default=False, is_active=True)
+ assert _can_access_pipeline(p, "admin", admin=True) is True
From 1203a4b75fa56e80a8781700763ea3cb94943895 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Sat, 7 Mar 2026 10:39:19 +0000
Subject: [PATCH 3/7] refactor(pipelines): address code review - shared
get_current_user_id, aria-live, deduplicate user ID logic
- Extract _get_user_id into shared auth.get_current_user_id() used by both
pipelines API and the assign-pipeline endpoint in files API
- Fix aria-live attribute: use two separate static containers (polite/assertive)
instead of dynamic Alpine.js binding for correct screen reader announcements
- Fix migration comment to accurately describe batch-mode FK creation
- Remove redundant tags parameter from reorder endpoint decorator
- Rename _make_file test helper to _make_test_file_record for clarity
- Update docs/UserGuide.md and docs/API.md with full Pipelines reference
Co-authored-by: christianlouis <361235+christianlouis@users.noreply.github.com>
---
app/api/files.py | 8 +-
app/api/pipelines.py | 15 +-
app/auth.py | 21 +++
docs/API.md | 170 +++++++++++++++++++++++
docs/UserGuide.md | 49 +++++++
frontend/templates/pipelines.html | 24 ++--
migrations/versions/017_add_pipelines.py | 5 +-
tests/test_api_pipelines.py | 8 +-
8 files changed, 271 insertions(+), 29 deletions(-)
diff --git a/app/api/files.py b/app/api/files.py
index ee5d2998..2277e2e9 100644
--- a/app/api/files.py
+++ b/app/api/files.py
@@ -1590,15 +1590,11 @@ def assign_pipeline_to_file(
HTTPException 404: If the file or pipeline does not exist / is not
accessible to the current user.
"""
- from app.auth import get_current_user
+ from app.auth import get_current_user, get_current_user_id
from app.models import Pipeline
user = get_current_user(request)
- # Derive user identity the same way the pipelines API does (_get_user_id)
- if user:
- user_id: str = user.get("preferred_username") or user.get("email") or user.get("id") or "anonymous"
- else:
- user_id = "anonymous"
+ user_id: str = get_current_user_id(request)
is_admin_user = bool(user and user.get("is_admin"))
diff --git a/app/api/pipelines.py b/app/api/pipelines.py
index 13182a80..d1ee6c23 100644
--- a/app/api/pipelines.py
+++ b/app/api/pipelines.py
@@ -18,7 +18,7 @@ from fastapi import APIRouter, Body, Depends, HTTPException, Request, status
from pydantic import BaseModel, Field
from sqlalchemy.orm import Session
-from app.auth import get_current_user, require_login
+from app.auth import get_current_user, get_current_user_id, require_login
from app.database import get_db
from app.models import Pipeline, PipelineStep
@@ -90,11 +90,12 @@ MAX_NAME_LENGTH = 255
def _get_user_id(request: Request) -> str:
- """Return a stable user identifier from the session."""
- user = get_current_user(request)
- if user:
- return user.get("preferred_username") or user.get("email") or user.get("id", "anonymous")
- return "anonymous"
+ """Return a stable user identifier from the session.
+
+ Delegates to :func:`app.auth.get_current_user_id` so the same fallback
+ logic ("anonymous") is used consistently throughout the application.
+ """
+ return get_current_user_id(request)
def _is_admin(request: Request) -> bool:
@@ -591,7 +592,7 @@ def add_step(pipeline_id: int, request: Request, db: DbSession, body: PipelineSt
return _serialize_step(step)
-@router.put("/{pipeline_id}/steps/reorder", tags=["pipelines"])
+@router.put("/{pipeline_id}/steps/reorder")
@require_login
def reorder_steps(
pipeline_id: int,
diff --git a/app/auth.py b/app/auth.py
index 3f8a43fa..c9c7099c 100644
--- a/app/auth.py
+++ b/app/auth.py
@@ -43,6 +43,27 @@ def get_current_user(request: Request):
return request.session.get("user")
+def get_current_user_id(request: Request) -> str:
+ """Return a stable string identifier for the authenticated user.
+
+ Falls back to ``"anonymous"`` when no user is in the session (e.g. when
+ AUTH_ENABLED=False in single-user mode). The returned value is consistent
+ with how pipeline and file ownership is stored in the database.
+
+ Priority order: ``preferred_username`` → ``email`` → ``id`` → ``"anonymous"``.
+
+ Args:
+ request: The current FastAPI request.
+
+ Returns:
+ A non-empty string identifying the current user.
+ """
+ user = get_current_user(request)
+ if not user or not isinstance(user, dict):
+ return "anonymous"
+ return user.get("preferred_username") or user.get("email") or user.get("id") or "anonymous"
+
+
def require_login(func):
if not AUTH_ENABLED:
return func # no-op
diff --git a/docs/API.md b/docs/API.md
index e736c260..2d016ddc 100644
--- a/docs/API.md
+++ b/docs/API.md
@@ -1441,6 +1441,176 @@ Execute a full data migration from source to target database.
```
+## Pipelines
+
+The pipeline API lets you build and manage custom document processing workflows. Each pipeline is owned by a single user (or by the system when `owner_id` is `null`).
+
+### Step-types catalogue
+
+```bash
+GET /api/pipelines/step-types
+```
+
+Returns the catalogue of built-in step types.
+
+**Response (200):**
+```json
+{
+ "convert_to_pdf": {
+ "label": "Convert to PDF",
+ "description": "Convert non-PDF documents to PDF format using Gotenberg.",
+ "config_schema": {}
+ },
+ "ocr": {
+ "label": "OCR Processing",
+ "description": "Extract text using Azure Document Intelligence or local Tesseract.",
+ "config_schema": {
+ "force_cloud_ocr": { "type": "boolean", "default": false }
+ }
+ }
+}
+```
+
+### List pipelines
+
+```bash
+GET /api/pipelines
+```
+
+Returns pipelines visible to the current user (own + system pipelines). Admins see all pipelines.
+
+### Create pipeline
+
+```bash
+POST /api/pipelines
+Content-Type: application/json
+
+{
+ "name": "My Workflow",
+ "description": "Converts, OCRs, and stores documents.",
+ "is_default": false,
+ "is_active": true
+}
+```
+
+**Response (201):**
+```json
+{
+ "id": 1,
+ "owner_id": "alice",
+ "name": "My Workflow",
+ "description": "Converts, OCRs, and stores documents.",
+ "is_default": false,
+ "is_active": true,
+ "created_at": "2026-03-07T10:00:00+00:00",
+ "updated_at": "2026-03-07T10:00:00+00:00"
+}
+```
+
+### Create system pipeline (admin only)
+
+```bash
+POST /api/pipelines/admin/system
+Content-Type: application/json
+
+{
+ "name": "Global Default",
+ "is_default": true
+}
+```
+
+### Get pipeline with steps
+
+```bash
+GET /api/pipelines/{pipeline_id}
+```
+
+**Response (200):**
+```json
+{
+ "id": 1,
+ "owner_id": "alice",
+ "name": "My Workflow",
+ "steps": [
+ { "id": 1, "position": 0, "step_type": "convert_to_pdf", "enabled": true, "config": {} },
+ { "id": 2, "position": 1, "step_type": "ocr", "enabled": true, "config": { "force_cloud_ocr": false } }
+ ]
+}
+```
+
+### Update pipeline
+
+```bash
+PUT /api/pipelines/{pipeline_id}
+Content-Type: application/json
+
+{ "name": "Renamed Workflow", "is_default": true }
+```
+
+### Delete pipeline
+
+```bash
+DELETE /api/pipelines/{pipeline_id}
+```
+
+Returns **204 No Content**.
+
+### Add step
+
+```bash
+POST /api/pipelines/{pipeline_id}/steps
+Content-Type: application/json
+
+{
+ "step_type": "ocr",
+ "label": "Cloud OCR",
+ "config": { "force_cloud_ocr": true },
+ "enabled": true
+}
+```
+
+### Update step
+
+```bash
+PUT /api/pipelines/{pipeline_id}/steps/{step_id}
+Content-Type: application/json
+
+{ "enabled": false }
+```
+
+### Delete step
+
+```bash
+DELETE /api/pipelines/{pipeline_id}/steps/{step_id}
+```
+
+Returns **204 No Content**.
+
+### Reorder steps
+
+```bash
+PUT /api/pipelines/{pipeline_id}/steps/reorder
+Content-Type: application/json
+
+[3, 1, 2]
+```
+
+Provide a complete ordered list of **all** step IDs. Their positions are reassigned 0, 1, 2, … in the given order.
+
+### Assign pipeline to a file
+
+```bash
+POST /api/files/{file_id}/assign-pipeline?pipeline_id=2
+```
+
+Pass no `pipeline_id` query parameter (or omit it) to clear the assignment.
+
+**Response (200):**
+```json
+{ "file_id": 42, "pipeline_id": 2 }
+```
+
+
## Further Assistance
For additional help with the API, please contact our support team or refer to the [Development Guide](../CONTRIBUTING.md).
diff --git a/docs/UserGuide.md b/docs/UserGuide.md
index 98ff365c..399791ea 100644
--- a/docs/UserGuide.md
+++ b/docs/UserGuide.md
@@ -334,6 +334,55 @@ PAPERLESS_CUSTOM_FIELDS_MAPPING='{"absender": "Sender", "empfaenger": "Recipient
3. After successful upload, custom fields are automatically populated
4. You can view the populated fields in your Paperless-ngx document details
+## Processing Pipelines
+
+Processing pipelines let you define exactly what happens to your documents when they are uploaded. Each pipeline is an ordered sequence of **steps** — for example: convert to PDF → OCR → extract metadata → send to storage.
+
+### Key concepts
+
+| Term | Meaning |
+|------|---------|
+| **Pipeline** | A named, ordered list of processing steps |
+| **Step** | A single processing action (e.g., OCR, metadata extraction) |
+| **System pipeline** | Created by an admin; visible to all users as a shared default |
+| **User pipeline** | Created by a regular user; private to that user |
+| **Default pipeline** | Marked `is_default=true`; used automatically for new uploads |
+
+### Managing your pipelines
+
+1. Navigate to **Pipelines** in the top navigation bar.
+2. Click **New Pipeline** to create a pipeline, give it a name and optional description.
+3. Expand the pipeline card and click **Add Step** to build the workflow.
+4. Use the ↑ / ↓ arrows to reorder steps, or click the edit icon to change step settings.
+5. Mark a pipeline as **Default** so new documents are automatically processed by it.
+
+### Available step types
+
+| Step Type | Description |
+|-----------|-------------|
+| `convert_to_pdf` | Convert non-PDF files to PDF using Gotenberg |
+| `check_duplicates` | Detect duplicate files by content hash |
+| `ocr` | Extract text with Azure Document Intelligence or local Tesseract |
+| `extract_metadata` | Extract structured metadata (type, sender, tags) with AI |
+| `embed_metadata` | Write extracted metadata into the PDF document properties |
+| `compute_embedding` | Compute semantic embeddings for similarity search |
+| `send_to_destinations` | Upload the processed document to all configured storage destinations |
+| `classify` | Classify the document type with AI |
+
+### Assigning a pipeline to a file
+
+You can assign (or change) the pipeline for an individual document via the file detail page or the API:
+
+```bash
+POST /api/files/{file_id}/assign-pipeline?pipeline_id=3
+```
+
+Pass no `pipeline_id` to clear the assignment and fall back to the system default.
+
+### Admin: system-wide pipelines
+
+Admins can create **system pipelines** that appear in every user's pipeline list. These can be set as the global default so all users benefit from a consistent processing baseline. Navigate to **Pipelines** and check the **System pipeline** box when creating a new one (admin only).
+
## API Access
For programmatic access, DocuElevate provides a comprehensive REST API:
diff --git a/frontend/templates/pipelines.html b/frontend/templates/pipelines.html
index fb36ec47..92d04b07 100644
--- a/frontend/templates/pipelines.html
+++ b/frontend/templates/pipelines.html
@@ -29,15 +29,21 @@
-
-
-
-
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/migrations/versions/017_add_pipelines.py b/migrations/versions/017_add_pipelines.py
index d14abbf9..b34f25c6 100644
--- a/migrations/versions/017_add_pipelines.py
+++ b/migrations/versions/017_add_pipelines.py
@@ -44,9 +44,8 @@ def upgrade() -> None:
sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now()),
)
- # Use batch mode for SQLite compatibility when altering the files table.
- # SQLite does not support adding FK constraints inline via ALTER TABLE, so we
- # add the plain integer column and define the FK reference at the model level.
+ # Use batch mode (copy-and-move strategy) for SQLite compatibility.
+ # A named FK constraint is created so Alembic can reference it in the downgrade.
with op.batch_alter_table("files") as batch_op:
batch_op.add_column(
sa.Column("pipeline_id", sa.Integer(), nullable=True),
diff --git a/tests/test_api_pipelines.py b/tests/test_api_pipelines.py
index ad7fcf33..f5daf532 100644
--- a/tests/test_api_pipelines.py
+++ b/tests/test_api_pipelines.py
@@ -16,7 +16,7 @@ from app.models import FileRecord, Pipeline, PipelineStep
# ---------------------------------------------------------------------------
-def _make_file(db_session, owner_id=None):
+def _make_test_file_record(db_session, owner_id=None):
"""Insert a minimal FileRecord and return it.
Default owner_id=None so tests work without an authenticated session.
@@ -344,7 +344,7 @@ class TestAssignPipelineToFile:
def test_assign_pipeline_to_file(self, client, db_session):
"""Assigning a pipeline to a file stores pipeline_id on the record."""
# File with no owner so anonymous test session can access it
- fr = _make_file(db_session, owner_id=None)
+ fr = _make_test_file_record(db_session, owner_id=None)
pipeline = client.post("/api/pipelines", json={"name": "Assign Test"}).json()
r = client.post(f"/api/files/{fr.id}/assign-pipeline?pipeline_id={pipeline['id']}")
@@ -358,7 +358,7 @@ class TestAssignPipelineToFile:
def test_clear_pipeline_from_file(self, client, db_session):
"""Passing no pipeline_id clears the assignment."""
- fr = _make_file(db_session, owner_id=None)
+ fr = _make_test_file_record(db_session, owner_id=None)
pipeline = client.post("/api/pipelines", json={"name": "Clearable"}).json()
client.post(f"/api/files/{fr.id}/assign-pipeline?pipeline_id={pipeline['id']}")
@@ -368,7 +368,7 @@ class TestAssignPipelineToFile:
def test_assign_nonexistent_pipeline_returns_404(self, client, db_session):
"""Assigning a non-existent pipeline returns 404."""
- fr = _make_file(db_session, owner_id=None)
+ fr = _make_test_file_record(db_session, owner_id=None)
r = client.post(f"/api/files/{fr.id}/assign-pipeline?pipeline_id=99999")
assert r.status_code == 404
From d318110bbe133e75a9cc3852d86295e6f682e2a4 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Sat, 7 Mar 2026 10:58:39 +0000
Subject: [PATCH 4/7] fix(pipelines): seed standard processing pipeline as
system default on startup
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
The pipeline management UI showed an empty list after first boot because no
default system pipeline was created. This adds seed_default_pipeline() which:
- Creates a system-owned (owner_id=NULL), is_default=True pipeline named
"Standard Processing Pipeline" at application startup
- Steps mirror the current hardcoded Celery processing workflow:
convert_to_pdf → check_duplicates → ocr → extract_metadata →
embed_metadata → compute_embedding → send_to_destinations
- Is idempotent: no-op if any system pipeline already exists
- Handles missing pipelines table gracefully (during first migration run)
Also wires the seeder into app/main.py lifespan startup using the same
pattern as seed_default_plans.
9 new tests added covering creation, step order, idempotency, and API visibility.
Co-authored-by: christianlouis <361235+christianlouis@users.noreply.github.com>
---
app/api/pipelines.py | 92 ++++++++++++++++++++++++++++
app/main.py | 14 +++++
tests/test_api_pipelines.py | 117 ++++++++++++++++++++++++++++++++++++
3 files changed, 223 insertions(+)
diff --git a/app/api/pipelines.py b/app/api/pipelines.py
index d1ee6c23..fc7bf851 100644
--- a/app/api/pipelines.py
+++ b/app/api/pipelines.py
@@ -792,3 +792,95 @@ def _unset_default(db: Session, owner_id: str | None) -> None:
db.query(Pipeline).filter(Pipeline.owner_id == owner_id, Pipeline.is_default.is_(True)).update(
{"is_default": False}
)
+
+
+# ---------------------------------------------------------------------------
+# Default system pipeline seeding
+# ---------------------------------------------------------------------------
+
+# The steps that make up the standard document-processing workflow. The order
+# here mirrors what the existing Celery-based pipeline executes for every
+# uploaded file.
+_DEFAULT_PIPELINE_STEPS: list[tuple[str, str]] = [
+ ("convert_to_pdf", "Convert to PDF"),
+ ("check_duplicates", "Check for Duplicates"),
+ ("ocr", "OCR Processing"),
+ ("extract_metadata", "Extract Metadata"),
+ ("embed_metadata", "Embed Metadata into PDF"),
+ ("compute_embedding", "Compute Text Embedding"),
+ ("send_to_destinations", "Send to Storage Destinations"),
+]
+
+#: Human-readable name shown in the management UI for the auto-seeded pipeline.
+DEFAULT_PIPELINE_NAME = "Standard Processing Pipeline"
+
+
+def seed_default_pipeline(db: Session) -> int:
+ """Ensure a system-owned default pipeline exists in the database.
+
+ This function is idempotent — it is a no-op when any system pipeline
+ (``owner_id IS NULL``) already exists. It is intended to be called once
+ at application startup (in ``app.main.lifespan``) so that the pipeline
+ management UI always shows the default workflow that mirrors the existing
+ Celery-based processing steps.
+
+ The created pipeline:
+
+ * ``owner_id = None`` — owned by the system, visible to all users
+ * ``is_default = True`` — selected automatically for new documents
+ * Steps (in order): convert_to_pdf → check_duplicates → ocr →
+ extract_metadata → embed_metadata → compute_embedding →
+ send_to_destinations
+
+ Args:
+ db: An active SQLAlchemy session.
+
+ Returns:
+ ``1`` if a new pipeline was created, ``0`` if one already existed.
+ """
+ try:
+ if db.query(Pipeline).filter(Pipeline.owner_id.is_(None)).count() > 0:
+ return 0
+ except Exception:
+ # Table may not exist yet during the very first migration run.
+ return 0
+
+ pipeline = Pipeline(
+ owner_id=None,
+ name=DEFAULT_PIPELINE_NAME,
+ description=(
+ "The standard document processing workflow: PDF conversion, "
+ "duplicate detection, OCR, metadata extraction and embedding, "
+ "semantic embeddings, and final distribution to storage destinations."
+ ),
+ is_default=True,
+ is_active=True,
+ )
+ db.add(pipeline)
+ try:
+ db.flush() # Assign pipeline.id without committing yet
+ except Exception as exc: # pragma: no cover
+ db.rollback()
+ logger.error(f"Failed to create default pipeline: {exc}")
+ return 0
+
+ for pos, (step_type, label) in enumerate(_DEFAULT_PIPELINE_STEPS):
+ db.add(
+ PipelineStep(
+ pipeline_id=pipeline.id,
+ position=pos,
+ step_type=step_type,
+ label=label,
+ enabled=True,
+ )
+ )
+
+ try:
+ db.commit()
+ logger.info("Seeded default system pipeline: '%s' (id=%d)", DEFAULT_PIPELINE_NAME, pipeline.id)
+ except Exception as exc: # pragma: no cover
+ db.rollback()
+ logger.error(f"Failed to seed default pipeline steps: {exc}")
+ return 0
+
+ return 1
diff --git a/app/main.py b/app/main.py
index 5f39ae2c..7c46593f 100644
--- a/app/main.py
+++ b/app/main.py
@@ -112,6 +112,20 @@ async def lifespan(app: FastAPI):
except Exception:
logging.debug("Subscription plan seeding skipped — DB may not be ready yet") # noqa: S110
+ # Seed the default system pipeline (mirrors the current hardcoded processing
+ # workflow) so it is immediately visible in the Pipelines management UI.
+ try:
+ from app.api.pipelines import seed_default_pipeline as _seed_pipeline
+ from app.database import SessionLocal as _SessionLocal # noqa: F811 (re-import for clarity)
+
+ _db_pipeline = _SessionLocal()
+ try:
+ _seed_pipeline(_db_pipeline)
+ finally:
+ _db_pipeline.close()
+ except Exception:
+ logging.debug("Default pipeline seeding skipped — DB may not be ready yet") # noqa: S110
+
# Application is now running
yield
diff --git a/tests/test_api_pipelines.py b/tests/test_api_pipelines.py
index f5daf532..a0d7785f 100644
--- a/tests/test_api_pipelines.py
+++ b/tests/test_api_pipelines.py
@@ -465,3 +465,120 @@ class TestPipelineAPIHelpers:
p = Pipeline(owner_id="user99", name="Private", is_default=False, is_active=True)
assert _can_access_pipeline(p, "admin", admin=True) is True
+
+
+# ---------------------------------------------------------------------------
+# Unit + integration tests – seed_default_pipeline
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.unit
+class TestSeedDefaultPipeline:
+ """Tests for the seed_default_pipeline startup helper."""
+
+ def test_seed_creates_pipeline(self, db_session):
+ """seed_default_pipeline creates exactly one system pipeline."""
+ from app.api.pipelines import seed_default_pipeline
+
+ result = seed_default_pipeline(db_session)
+
+ assert result == 1
+ pipelines = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).all()
+ assert len(pipelines) == 1
+
+ def test_seeded_pipeline_is_default(self, db_session):
+ """The seeded pipeline has is_default=True and is_active=True."""
+ from app.api.pipelines import seed_default_pipeline
+
+ seed_default_pipeline(db_session)
+ p = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).first()
+
+ assert p is not None
+ assert p.is_default is True
+ assert p.is_active is True
+
+ def test_seeded_pipeline_has_correct_name(self, db_session):
+ """The seeded pipeline uses the canonical DEFAULT_PIPELINE_NAME."""
+ from app.api.pipelines import DEFAULT_PIPELINE_NAME, seed_default_pipeline
+
+ seed_default_pipeline(db_session)
+ p = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).first()
+
+ assert p.name == DEFAULT_PIPELINE_NAME
+
+ def test_seeded_pipeline_steps_count(self, db_session):
+ """The seeded pipeline has the correct number of steps."""
+ from app.api.pipelines import _DEFAULT_PIPELINE_STEPS, seed_default_pipeline
+
+ seed_default_pipeline(db_session)
+ p = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).first()
+ steps = db_session.query(PipelineStep).filter(PipelineStep.pipeline_id == p.id).all()
+
+ assert len(steps) == len(_DEFAULT_PIPELINE_STEPS)
+
+ def test_seeded_pipeline_step_types_and_order(self, db_session):
+ """Steps are in the correct order and match the expected step types."""
+ from app.api.pipelines import _DEFAULT_PIPELINE_STEPS, seed_default_pipeline
+
+ seed_default_pipeline(db_session)
+ p = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).first()
+ steps = (
+ db_session.query(PipelineStep)
+ .filter(PipelineStep.pipeline_id == p.id)
+ .order_by(PipelineStep.position)
+ .all()
+ )
+
+ expected_types = [step_type for step_type, _ in _DEFAULT_PIPELINE_STEPS]
+ actual_types = [s.step_type for s in steps]
+ assert actual_types == expected_types
+
+ def test_seeded_pipeline_all_steps_enabled(self, db_session):
+ """All seeded steps are enabled by default."""
+ from app.api.pipelines import seed_default_pipeline
+
+ seed_default_pipeline(db_session)
+ p = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).first()
+ steps = db_session.query(PipelineStep).filter(PipelineStep.pipeline_id == p.id).all()
+
+ assert all(s.enabled for s in steps), "All seeded steps should be enabled"
+
+ def test_seed_is_idempotent(self, db_session):
+ """Calling seed_default_pipeline twice does not create a duplicate."""
+ from app.api.pipelines import seed_default_pipeline
+
+ first = seed_default_pipeline(db_session)
+ second = seed_default_pipeline(db_session)
+
+ assert first == 1
+ assert second == 0 # No-op on second call
+
+ count = db_session.query(Pipeline).filter(Pipeline.owner_id.is_(None)).count()
+ assert count == 1
+
+ def test_seeded_pipeline_visible_via_api(self, client):
+ """The default pipeline is visible in the GET /api/pipelines listing."""
+ from app.api.pipelines import DEFAULT_PIPELINE_NAME, seed_default_pipeline
+ from app.database import get_db
+
+ # Seed using the same DB session that the test client uses
+ db = next(client.app.dependency_overrides[get_db]())
+ seed_default_pipeline(db)
+
+ r = client.get("/api/pipelines")
+ assert r.status_code == 200
+ names = [p["name"] for p in r.json()]
+ assert DEFAULT_PIPELINE_NAME in names
+
+ def test_seeded_pipeline_default_flag_visible_via_api(self, client):
+ """The seeded pipeline is returned with is_default=True via the API."""
+ from app.api.pipelines import DEFAULT_PIPELINE_NAME, seed_default_pipeline
+ from app.database import get_db
+
+ db = next(client.app.dependency_overrides[get_db]())
+ seed_default_pipeline(db)
+
+ r = client.get("/api/pipelines")
+ default_pipelines = [p for p in r.json() if p["name"] == DEFAULT_PIPELINE_NAME]
+ assert len(default_pipelines) == 1
+ assert default_pipelines[0]["is_default"] is True
From a644efe016220af088654e6448b9429881bcc27b Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Sat, 7 Mar 2026 11:14:49 +0000
Subject: [PATCH 5/7] feat(files): show assigned pipeline info on file status
and detail views
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
The file detail page (/files/{id}/detail) and file view page (/files/{id})
previously showed no information about which processing pipeline was used.
Changes:
- _STEP_TYPE_TO_STAGES mapping: pipeline step_type → Celery log stage keys
(with maintenance comment requiring updates when new step types are added)
- _ALWAYS_SHOW_STAGES: stages always visible regardless of pipeline
- _resolve_pipeline(db, file_record): resolves the pipeline for a file —
uses explicit pipeline_id when set, falls back to active system default
- _compute_processing_flow: new pipeline_steps parameter; when provided,
filters flow graph to only show stages for the pipeline's enabled steps
(+ always-show stages + any stage that actually ran). Also adds
convert_to_pdf to the flow stage catalogue.
- file_detail_page: passes pipeline_info + pipeline-filtered flow_data
- file_view_page: passes pipeline_info
Templates:
- file_detail.html: 'Processing Pipeline' detail row with name link and
colour-coded badge (System Default / System / Custom)
- file_view.html: 'Pipeline' info row in sidebar with (default)/(custom) tag
Tests:
- TestPipelineInfoInViews with 14 tests covering _resolve_pipeline,
_compute_processing_flow filtering, completeness assertion for
_STEP_TYPE_TO_STAGES, and HTTP-level view tests
Co-authored-by: christianlouis <361235+christianlouis@users.noreply.github.com>
---
app/views/files.py | 116 +++++++++-
frontend/templates/file_detail.html | 17 ++
frontend/templates/file_view.html | 15 ++
tests/test_views_files_comprehensive.py | 268 ++++++++++++++++++++++++
4 files changed, 412 insertions(+), 4 deletions(-)
diff --git a/app/views/files.py b/app/views/files.py
index c1a2c7da..58efa132 100644
--- a/app/views/files.py
+++ b/app/views/files.py
@@ -270,6 +270,9 @@ def file_view_page(request: Request, file_id: int, db: Session = Depends(get_db)
except Exception:
step_summary = None
+ # Resolve the pipeline assigned to this file (explicit or system default)
+ pipeline_info = _resolve_pipeline(db, file_record)
+
return templates.TemplateResponse(
"file_view.html",
{
@@ -279,6 +282,7 @@ def file_view_page(request: Request, file_id: int, db: Session = Depends(get_db)
"original_file_exists": original_file_exists,
"processed_file_exists": processed_file_exists,
"step_summary": step_summary,
+ "pipeline_info": pipeline_info,
},
)
except Exception as e:
@@ -337,8 +341,11 @@ def file_detail_page(request: Request, file_id: int, db: Session = Depends(get_d
except Exception as e:
logger.warning(f"Failed to load metadata from {metadata_path}: {e}")
- # Compute processing flow for visualization
- flow_data = _compute_processing_flow(logs)
+ # Resolve the pipeline assigned to this file (explicit or system default)
+ pipeline_info = _resolve_pipeline(db, file_record)
+
+ # Compute processing flow for visualization — filter to pipeline steps when available
+ flow_data = _compute_processing_flow(logs, pipeline_steps=pipeline_info["steps"] if pipeline_info else None)
# Compute step-aligned summary from status table (preferred) or fallback to logs
try:
@@ -360,6 +367,7 @@ def file_detail_page(request: Request, file_id: int, db: Session = Depends(get_d
"gpt_metadata": gpt_metadata,
"flow_data": flow_data,
"step_summary": step_summary,
+ "pipeline_info": pipeline_info,
},
)
except Exception as e:
@@ -367,15 +375,99 @@ def file_detail_page(request: Request, file_id: int, db: Session = Depends(get_d
return templates.TemplateResponse("file_detail.html", {"request": request, "file": None, "error": str(e)})
-def _compute_processing_flow(logs):
+# ---------------------------------------------------------------------------
+# Pipeline ↔ Celery-log stage mapping
+# ---------------------------------------------------------------------------
+
+# Maps each pipeline step_type to the set of Celery task log stage keys that
+# implement it. Used to filter the flow visualization when a pipeline is
+# assigned to a file.
+#
+# ⚠️ MAINTENANCE NOTE: When a new step type is added to PIPELINE_STEP_TYPES
+# in app/api/pipelines.py it MUST also be added here, otherwise the flow
+# visualization will silently skip its Celery-task stages for files using that
+# step type. The test ``TestPipelineInfoInViews::test_step_type_mapping_is_complete``
+# enforces this invariant automatically.
+_STEP_TYPE_TO_STAGES: dict[str, list[str]] = {
+ "convert_to_pdf": ["convert_to_pdf"],
+ "check_duplicates": ["check_for_duplicates"],
+ "ocr": ["check_text", "extract_text", "process_with_ocr"],
+ "extract_metadata": ["extract_metadata_with_gpt"],
+ "embed_metadata": ["embed_metadata_into_pdf"],
+ "compute_embedding": ["compute_embedding"],
+ "send_to_destinations": ["finalize_document_storage", "send_to_all_destinations"],
+ # "classify" is defined in PIPELINE_STEP_TYPES but has no Celery log stages yet.
+ # When a classify task is implemented, add its stage key(s) here.
+ "classify": [],
+}
+
+# These internal bookkeeping stages are always shown in the flow regardless of
+# which pipeline steps are defined.
+_ALWAYS_SHOW_STAGES: frozenset[str] = frozenset({"create_file_record"})
+
+
+def _resolve_pipeline(db: Session, file_record) -> dict | None:
+ """Resolve the pipeline information for a file.
+
+ If the file has an explicit ``pipeline_id``, load that pipeline.
+ Otherwise fall back to the active system-default pipeline
+ (``owner_id IS NULL``, ``is_default=True``).
+
+ Returns a dict with keys:
+ id, name, description, is_default, is_system, is_explicit, steps
+ or ``None`` when no pipeline exists in the database.
+ """
+ from app.models import Pipeline, PipelineStep
+
+ pipeline = None
+ if file_record.pipeline_id:
+ pipeline = db.query(Pipeline).filter(Pipeline.id == file_record.pipeline_id).first()
+
+ if pipeline is None:
+ pipeline = (
+ db.query(Pipeline)
+ .filter(
+ Pipeline.owner_id.is_(None),
+ Pipeline.is_default.is_(True),
+ Pipeline.is_active.is_(True),
+ )
+ .first()
+ )
+
+ if pipeline is None:
+ return None
+
+ steps = db.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline.id).order_by(PipelineStep.position).all()
+
+ return {
+ "id": pipeline.id,
+ "name": pipeline.name,
+ "description": pipeline.description,
+ "is_default": pipeline.is_default,
+ "is_system": pipeline.owner_id is None,
+ # True when the file has a pipeline explicitly assigned (not inferred default)
+ "is_explicit": bool(file_record.pipeline_id),
+ "steps": steps,
+ }
+
+
+def _compute_processing_flow(logs, pipeline_steps=None):
"""
Compute the processing flow structure from logs for visualization.
Returns a structured representation of the processing pipeline with branches.
Detects upload sub-tasks and organizes them as branches under the parent upload stage.
+
+ Args:
+ logs: list of ProcessingLog objects (ordered by timestamp asc)
+ pipeline_steps: optional list of PipelineStep objects for the assigned pipeline.
+ When provided, the set of stages shown is filtered to only those that
+ correspond to the pipeline's enabled steps (plus bookkeeping stages like
+ ``create_file_record`` and any stage that actually ran in the logs).
"""
- # Define the main processing stages
+ # Define the full catalogue of main processing stages
stages = {
+ "convert_to_pdf": {"label": "Convert to PDF", "next": ["check_for_duplicates", "create_file_record"]},
"check_for_duplicates": {"label": "Check for Duplicates", "next": ["create_file_record"]},
"create_file_record": {"label": "Create File Record", "next": ["check_text"]},
"check_text": {
@@ -406,6 +498,22 @@ def _compute_processing_flow(logs):
if "create_file_record" in stages:
stages["create_file_record"]["next"] = ["check_text"]
+ # When a pipeline is assigned, filter stages to only those relevant to the
+ # pipeline's enabled steps plus always-show bookkeeping stages and any stage
+ # that actually produced log entries (so nothing already-run is hidden).
+ if pipeline_steps is not None:
+ # Collect Celery stage keys that the pipeline's enabled steps map to
+ allowed: set[str] = set(_ALWAYS_SHOW_STAGES)
+ for ps in pipeline_steps:
+ if ps.enabled:
+ allowed.update(_STEP_TYPE_TO_STAGES.get(ps.step_type, []))
+ # Pre-scan logs so we can also keep any stage that already ran
+ ran_stages: set[str] = set()
+ for log in logs:
+ ran_stages.add(log.step_name)
+ allowed.update(ran_stages)
+ stages = {k: v for k, v in stages.items() if k in allowed}
+
# Define upload sub-tasks (branches)
upload_tasks = {
"upload_to_dropbox": "Dropbox",
diff --git a/frontend/templates/file_detail.html b/frontend/templates/file_detail.html
index c9d34f57..70b150be 100644
--- a/frontend/templates/file_detail.html
+++ b/frontend/templates/file_detail.html
@@ -1192,6 +1192,23 @@
Created At{{ file.created_at.strftime('%Y-%m-%d %H:%M:%S') if file.created_at else 'N/A' }}
+
+ Processing Pipeline
+
+ {% if pipeline_info %}
+ {{ pipeline_info.name }}
+ {% if not pipeline_info.is_explicit and pipeline_info.is_system %}
+ System Default
+ {% elif pipeline_info.is_system %}
+ System
+ {% else %}
+ Custom
+ {% endif %}
+ {% else %}
+ Standard (legacy)
+ {% endif %}
+
+