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 %}
+
+
Original File Path
diff --git a/frontend/templates/file_view.html b/frontend/templates/file_view.html
index 873ab6e7..cf5a4024 100644
--- a/frontend/templates/file_view.html
+++ b/frontend/templates/file_view.html
@@ -301,6 +301,21 @@
Hash
{{ file.filehash[:32] }}…
+
+
Pipeline
+
+ {% if pipeline_info %}
+ {{ pipeline_info.name }}
+ {% if not pipeline_info.is_explicit %}
+ (default)
+ {% elif not pipeline_info.is_system %}
+ (custom)
+ {% endif %}
+ {% else %}
+ Standard
+ {% endif %}
+
+
diff --git a/tests/test_views_files_comprehensive.py b/tests/test_views_files_comprehensive.py
index d77f2eb4..d522a637 100644
--- a/tests/test_views_files_comprehensive.py
+++ b/tests/test_views_files_comprehensive.py
@@ -1713,3 +1713,271 @@ class TestGetTextWithContent:
assert "text" in data
assert data["text"] # Should have non-empty text
assert "No text" not in data["text"]
+
+
+# ---------------------------------------------------------------------------
+# Pipeline-info tests: file_detail and file_view views
+# ---------------------------------------------------------------------------
+
+
+@pytest.mark.unit
+class TestPipelineInfoInViews:
+ """Tests that pipeline information is correctly resolved and passed to templates."""
+
+ # ------------------------------------------------------------------
+ # Helpers
+ # ------------------------------------------------------------------
+
+ @staticmethod
+ def _make_file(db_session, pipeline_id=None):
+ from app.models import FileRecord
+
+ f = FileRecord(
+ filehash="ph_" + str(pipeline_id),
+ original_filename="doc.pdf",
+ local_filename="/tmp/doc.pdf",
+ file_size=1024,
+ mime_type="application/pdf",
+ pipeline_id=pipeline_id,
+ )
+ db_session.add(f)
+ db_session.commit()
+ return f
+
+ @staticmethod
+ def _make_system_pipeline(db_session, is_default=True):
+ from app.models import Pipeline, PipelineStep
+
+ p = Pipeline(
+ owner_id=None,
+ name="Standard Processing Pipeline",
+ description="System default",
+ is_default=is_default,
+ is_active=True,
+ )
+ db_session.add(p)
+ db_session.flush()
+
+ for pos, (step_type, label) in enumerate(
+ [
+ ("convert_to_pdf", "Convert to PDF"),
+ ("ocr", "OCR"),
+ ("send_to_destinations", "Send"),
+ ]
+ ):
+ db_session.add(PipelineStep(pipeline_id=p.id, position=pos, step_type=step_type, label=label, enabled=True))
+
+ db_session.commit()
+ return p
+
+ @staticmethod
+ def _make_custom_pipeline(db_session):
+ from app.models import Pipeline, PipelineStep
+
+ p = Pipeline(
+ owner_id="user1",
+ name="My Custom Pipeline",
+ is_default=False,
+ is_active=True,
+ )
+ db_session.add(p)
+ db_session.flush()
+ db_session.add(PipelineStep(pipeline_id=p.id, position=0, step_type="ocr", label="OCR", enabled=True))
+ db_session.commit()
+ return p
+
+ # ------------------------------------------------------------------
+ # _resolve_pipeline unit tests
+ # ------------------------------------------------------------------
+
+ def test_resolve_pipeline_explicit_assignment(self, db_session):
+ """File with an explicit pipeline_id resolves to that pipeline."""
+ from app.views.files import _resolve_pipeline
+
+ pipeline = self._make_system_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=pipeline.id)
+
+ info = _resolve_pipeline(db_session, file_rec)
+
+ assert info is not None
+ assert info["id"] == pipeline.id
+ assert info["is_explicit"] is True
+ assert info["is_system"] is True
+
+ def test_resolve_pipeline_fallback_to_system_default(self, db_session):
+ """File without pipeline_id falls back to system-default pipeline."""
+ from app.views.files import _resolve_pipeline
+
+ pipeline = self._make_system_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=None)
+
+ info = _resolve_pipeline(db_session, file_rec)
+
+ assert info is not None
+ assert info["id"] == pipeline.id
+ assert info["is_explicit"] is False
+ assert info["is_system"] is True
+ assert info["is_default"] is True
+
+ def test_resolve_pipeline_returns_none_when_no_pipeline_in_db(self, db_session):
+ """Returns None when no pipeline exists (empty database)."""
+ from app.views.files import _resolve_pipeline
+
+ file_rec = self._make_file(db_session, pipeline_id=None)
+ info = _resolve_pipeline(db_session, file_rec)
+
+ assert info is None
+
+ def test_resolve_pipeline_includes_steps(self, db_session):
+ """Returned dict contains the pipeline's steps in order."""
+ from app.views.files import _resolve_pipeline
+
+ pipeline = self._make_system_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=pipeline.id)
+
+ info = _resolve_pipeline(db_session, file_rec)
+
+ assert info is not None
+ assert len(info["steps"]) == 3
+ assert info["steps"][0].step_type == "convert_to_pdf"
+
+ def test_resolve_pipeline_custom_pipeline(self, db_session):
+ """File with an explicit custom pipeline resolves correctly."""
+ from app.views.files import _resolve_pipeline
+
+ pipeline = self._make_custom_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=pipeline.id)
+
+ info = _resolve_pipeline(db_session, file_rec)
+
+ assert info is not None
+ assert info["id"] == pipeline.id
+ assert info["name"] == "My Custom Pipeline"
+ assert info["is_system"] is False
+ assert info["is_explicit"] is True
+
+ # ------------------------------------------------------------------
+ # _compute_processing_flow pipeline filtering
+ # ------------------------------------------------------------------
+
+ def test_compute_flow_without_pipeline_shows_all_stages(self):
+ """Without a pipeline, all hardcoded stages are included."""
+ from app.views.files import _compute_processing_flow
+
+ flow = _compute_processing_flow([], pipeline_steps=None)
+ # Should include standard stages (create_file_record, check_text, etc.)
+ keys = {s["key"] for s in flow}
+ assert "create_file_record" in keys
+ assert "extract_metadata_with_gpt" in keys
+
+ def test_step_type_mapping_is_complete(self):
+ """Every step type in PIPELINE_STEP_TYPES has an entry in _STEP_TYPE_TO_STAGES."""
+ from app.api.pipelines import PIPELINE_STEP_TYPES
+ from app.views.files import _STEP_TYPE_TO_STAGES
+
+ missing = set(PIPELINE_STEP_TYPES.keys()) - set(_STEP_TYPE_TO_STAGES.keys())
+ assert not missing, (
+ f"The following pipeline step types are missing from _STEP_TYPE_TO_STAGES "
+ f"in app/views/files.py: {missing}. "
+ "Add them with their corresponding Celery log stage key(s) (use [] if none yet)."
+ )
+
+ def test_compute_flow_with_pipeline_filters_stages(self, db_session):
+ """With a pipeline, only mapped stages are shown (plus always-show and ran stages)."""
+
+ from app.models import PipelineStep
+ from app.views.files import _compute_processing_flow
+
+ pipeline = self._make_system_pipeline(db_session)
+ # pipeline has: convert_to_pdf, ocr, send_to_destinations
+
+ # Query steps explicitly (no SQLAlchemy relationship defined on Pipeline)
+ steps = db_session.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline.id).all()
+
+ flow = _compute_processing_flow([], pipeline_steps=steps)
+ # When no logs and pipeline steps provided, only pipeline-mapped + always-show stages appear
+ keys = {s["key"] for s in flow}
+ # Always show
+ assert "create_file_record" in keys
+ # ocr maps to check_text / extract_text / process_with_ocr
+ assert "check_text" in keys or "extract_text" in keys
+ # embed_metadata not in pipeline → should be absent (no logs ran it)
+ assert "embed_metadata_into_pdf" not in keys
+
+ def test_compute_flow_with_pipeline_always_shows_ran_stages(self, db_session):
+ """Stages that actually ran are always shown even if not in the pipeline."""
+ from unittest.mock import Mock
+
+ from app.models import PipelineStep
+ from app.views.files import _compute_processing_flow
+
+ pipeline = self._make_system_pipeline(db_session)
+ steps = db_session.query(PipelineStep).filter(PipelineStep.pipeline_id == pipeline.id).all()
+
+ # Simulate a log entry for embed_metadata_into_pdf (not in this pipeline)
+ ran_log = Mock(
+ step_name="embed_metadata_into_pdf",
+ status="success",
+ message="Done",
+ timestamp=Mock(),
+ task_id="t1",
+ )
+
+ flow = _compute_processing_flow([ran_log], pipeline_steps=steps)
+ keys = {s["key"] for s in flow}
+ assert "embed_metadata_into_pdf" in keys
+
+ # ------------------------------------------------------------------
+ # Integration: view endpoints pass pipeline_info to template
+ # ------------------------------------------------------------------
+
+ def test_file_detail_page_includes_pipeline_name(self, client, db_session):
+ """GET /files/{id}/detail response body contains the pipeline name."""
+ pipeline = self._make_system_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=None)
+
+ response = client.get(f"/files/{file_rec.id}/detail")
+
+ assert response.status_code == 200
+ assert b"Standard Processing Pipeline" in response.content
+
+ def test_file_detail_page_shows_system_default_badge(self, client, db_session):
+ """File without explicit pipeline shows 'System Default' badge in detail view."""
+ self._make_system_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=None)
+
+ response = client.get(f"/files/{file_rec.id}/detail")
+
+ assert response.status_code == 200
+ assert b"System Default" in response.content
+
+ def test_file_detail_page_shows_custom_badge_for_custom_pipeline(self, client, db_session):
+ """File with a custom (non-system) pipeline shows 'Custom' badge."""
+ pipeline = self._make_custom_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=pipeline.id)
+
+ response = client.get(f"/files/{file_rec.id}/detail")
+
+ assert response.status_code == 200
+ assert b"My Custom Pipeline" in response.content
+ assert b"Custom" in response.content
+
+ def test_file_view_page_includes_pipeline_name(self, client, db_session):
+ """GET /files/{id} response body contains the pipeline name in the sidebar."""
+ pipeline = self._make_system_pipeline(db_session)
+ file_rec = self._make_file(db_session, pipeline_id=None)
+
+ response = client.get(f"/files/{file_rec.id}")
+
+ assert response.status_code == 200
+ assert b"Standard Processing Pipeline" in response.content
+
+ def test_file_view_page_no_pipeline_shows_standard(self, client, db_session):
+ """When no pipeline exists, file view shows 'Standard' fallback text."""
+ # No pipeline in DB
+ file_rec = self._make_file(db_session, pipeline_id=None)
+
+ response = client.get(f"/files/{file_rec.id}")
+
+ assert response.status_code == 200
+ assert b"Standard" in response.content