Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 12 additions & 1 deletion app/api/routes/__init__.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,21 @@
"""Route packages available for import convenience."""

from . import datasets, health, projects, repositories, sessions, stats, uploads, users
from . import (
datasets,
health,
pipeline,
projects,
repositories,
sessions,
stats,
uploads,
users,
)

__all__ = [
"datasets",
"health",
"pipeline",
"projects",
"repositories",
"sessions",
Expand Down
31 changes: 31 additions & 0 deletions app/api/routes/pipeline.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
"""Read-only endpoints exposing the canonical analysis pipeline definition."""

from fastapi import APIRouter

from ... import schemas
from ...core import get_default_pipeline

router = APIRouter(tags=["pipeline"])


@router.get(
"/pipeline",
response_model=schemas.PipelineDefinitionRead,
summary="Get canonical pipeline definition",
)
async def get_pipeline_definition() -> schemas.PipelineDefinitionRead:
"""Return the shared pipeline metadata including ordered steps."""

steps = [
schemas.PipelineDefinitionStep(
order=index,
name=step.name,
title=step.title,
description=step.description,
)
for index, step in enumerate(get_default_pipeline(), start=1)
]
return schemas.PipelineDefinitionRead(steps=steps, total_steps=len(steps))


__all__ = ["router"]
57 changes: 10 additions & 47 deletions app/api/routes/uploads.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
from datetime import datetime, timezone
from hashlib import sha256
from pathlib import Path
from typing import Any, Awaitable, Callable
from typing import Any

from fastapi import APIRouter, File, HTTPException, Response, UploadFile, status
from fastapi.responses import FileResponse
Expand All @@ -20,6 +20,7 @@

from ... import models, schemas
from ...config import get_settings
from ...core import PipelineStepDefinition, get_default_pipeline
from ...database import SessionLocal
from ..dependencies import DBSession

Expand Down Expand Up @@ -212,45 +213,7 @@ def _hash_file(path: Path) -> str:
return hasher.hexdigest()


async def _simulate_step_runtime(seconds: int) -> None:
"""Sleep in one-second increments to emulate long-running work."""

for _ in range(seconds):
await asyncio.sleep(1)


async def _step_ingest() -> None:
await _simulate_step_runtime(10)


async def _step_calibration() -> None:
await _simulate_step_runtime(10)


async def _step_registration() -> None:
await _simulate_step_runtime(10)


async def _step_photometry() -> None:
await _simulate_step_runtime(10)


async def _step_classification() -> None:
await _simulate_step_runtime(10)


async def _step_reporting() -> None:
await _simulate_step_runtime(10)


SIMULATED_PIPELINE_STEPS: list[tuple[str, Callable[[], Awaitable[None]]]] = [
("ingest", _step_ingest),
("calibration", _step_calibration),
("registration", _step_registration),
("photometry", _step_photometry),
("classification", _step_classification),
("reporting", _step_reporting),
]
SIMULATED_PIPELINE_STEPS: tuple[PipelineStepDefinition, ...] = get_default_pipeline()

_ACTIVE_SESSION_LOCK = threading.Lock()
_ACTIVE_SESSION_WORKERS: set[int] = set()
Expand Down Expand Up @@ -309,7 +272,7 @@ async def _process() -> None:
return

session_obj.status = "running"
session_obj.current_step = SIMULATED_PIPELINE_STEPS[0][0]
session_obj.current_step = SIMULATED_PIPELINE_STEPS[0].name
session_obj.started_at = session_obj.started_at or datetime.now(tz=timezone.utc)
session_obj.progress = 0
await session.commit()
Expand All @@ -321,12 +284,12 @@ async def _process() -> None:
existing_steps = {step.step_name: step for step in result}

pipeline_steps: list[models.PipelineStep] = []
for step_name, _ in SIMULATED_PIPELINE_STEPS:
step_record = existing_steps.get(step_name)
for step_definition in SIMULATED_PIPELINE_STEPS:
step_record = existing_steps.get(step_definition.name)
if step_record is None:
step_record = models.PipelineStep(
run_id=session_obj.run_id,
step_name=step_name,
step_name=step_definition.name,
status="queued",
progress=0,
)
Expand All @@ -342,18 +305,18 @@ async def _process() -> None:

total_steps = len(SIMULATED_PIPELINE_STEPS)

for index, (step_name, step_runner) in enumerate(SIMULATED_PIPELINE_STEPS):
for index, step_definition in enumerate(SIMULATED_PIPELINE_STEPS):
step_record = pipeline_steps[index]
step_record.status = "running"
step_record.started_at = datetime.now(tz=timezone.utc)
step_record.progress = 0

session_obj.current_step = step_name
session_obj.current_step = step_definition.name
session_obj.status = "running"
session_obj.progress = int((index / total_steps) * 100)
await session.commit()

await step_runner()
await step_definition.runner()

step_record.status = "completed"
step_record.finished_at = datetime.now(tz=timezone.utc)
Expand Down
17 changes: 17 additions & 0 deletions app/core/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
"""Core analysis step definitions used to assemble processing pipelines."""

from .pipeline import (
DEFAULT_PIPELINE,
PipelineStepDefinition,
get_default_pipeline,
get_step,
iter_step_names,
)

__all__ = [
"DEFAULT_PIPELINE",
"PipelineStepDefinition",
"get_default_pipeline",
"get_step",
"iter_step_names",
]
11 changes: 11 additions & 0 deletions app/core/calibration.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated calibration step applying dark/bias/flat corrections."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "calibration"
STEP_TITLE = "Calibrate exposures"
STEP_DESCRIPTION = "Apply dark, bias, and flat-field corrections to raw frames."

run_step = bind_runtime(2.0)
11 changes: 11 additions & 0 deletions app/core/classification.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated classification step for candidate events."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "classification"
STEP_TITLE = "Classify events"
STEP_DESCRIPTION = "Score extracted signals to prioritise follow-up."

run_step = bind_runtime(1.4)
11 changes: 11 additions & 0 deletions app/core/denoise.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated denoising step for extracted lightcurves."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "denoise"
STEP_TITLE = "Denoise lightcurve"
STEP_DESCRIPTION = "Reduce noise using frequency-domain filtering techniques."

run_step = bind_runtime(1.0)
11 changes: 11 additions & 0 deletions app/core/ingest.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated ingest step for the analysis pipeline."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "ingest"
STEP_TITLE = "Ingest raw exposures"
STEP_DESCRIPTION = "Pull uploaded FITS files into the processing workspace."

run_step = bind_runtime(1.5)
11 changes: 11 additions & 0 deletions app/core/lightcurve.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated lightcurve extraction step."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "lightcurve"
STEP_TITLE = "Extract lightcurve"
STEP_DESCRIPTION = "Generate time-series photometry for the aligned exposures."

run_step = bind_runtime(1.8)
103 changes: 103 additions & 0 deletions app/core/pipeline.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
"""Declarative pipeline configuration backed by core analysis steps."""

from __future__ import annotations

from dataclasses import dataclass
from typing import Awaitable, Callable, Iterable, Mapping

from . import (
calibration,
classification,
denoise,
ingest,
lightcurve,
registration,
reporting,
)


@dataclass(frozen=True)
class PipelineStepDefinition:
"""Metadata describing a single pipeline step."""

name: str
title: str
description: str
runner: Callable[[], Awaitable[None]]


DEFAULT_PIPELINE: tuple[PipelineStepDefinition, ...] = (
PipelineStepDefinition(
name=ingest.STEP_NAME,
title=ingest.STEP_TITLE,
description=ingest.STEP_DESCRIPTION,
runner=ingest.run_step,
),
PipelineStepDefinition(
name=calibration.STEP_NAME,
title=calibration.STEP_TITLE,
description=calibration.STEP_DESCRIPTION,
runner=calibration.run_step,
),
PipelineStepDefinition(
name=registration.STEP_NAME,
title=registration.STEP_TITLE,
description=registration.STEP_DESCRIPTION,
runner=registration.run_step,
),
PipelineStepDefinition(
name=lightcurve.STEP_NAME,
title=lightcurve.STEP_TITLE,
description=lightcurve.STEP_DESCRIPTION,
runner=lightcurve.run_step,
),
PipelineStepDefinition(
name=denoise.STEP_NAME,
title=denoise.STEP_TITLE,
description=denoise.STEP_DESCRIPTION,
runner=denoise.run_step,
),
PipelineStepDefinition(
name=classification.STEP_NAME,
title=classification.STEP_TITLE,
description=classification.STEP_DESCRIPTION,
runner=classification.run_step,
),
PipelineStepDefinition(
name=reporting.STEP_NAME,
title=reporting.STEP_TITLE,
description=reporting.STEP_DESCRIPTION,
runner=reporting.run_step,
),
)

_STEP_LOOKUP: Mapping[str, PipelineStepDefinition] = {
step.name: step for step in DEFAULT_PIPELINE
}


def get_default_pipeline() -> tuple[PipelineStepDefinition, ...]:
"""Return the immutable default pipeline definition."""

return DEFAULT_PIPELINE


def get_step(name: str) -> PipelineStepDefinition | None:
"""Return metadata for a named step if it exists."""

return _STEP_LOOKUP.get(name)


def iter_step_names() -> Iterable[str]:
"""Iterate over canonical step names in execution order."""

return (step.name for step in DEFAULT_PIPELINE)


__all__ = [
"PipelineStepDefinition",
"DEFAULT_PIPELINE",
"get_default_pipeline",
"get_step",
"iter_step_names",
]
11 changes: 11 additions & 0 deletions app/core/registration.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated registration step for aligning frames."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "registration"
STEP_TITLE = "Register exposures"
STEP_DESCRIPTION = "Align frames against reference stars to stabilise the stack."

run_step = bind_runtime(1.2)
11 changes: 11 additions & 0 deletions app/core/reporting.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
"""Simulated reporting step producing summary artefacts."""

from __future__ import annotations

from .runtime import bind_runtime

STEP_NAME = "reporting"
STEP_TITLE = "Compile report"
STEP_DESCRIPTION = "Generate session-level reports and notify subscribers."

run_step = bind_runtime(0.8)
27 changes: 27 additions & 0 deletions app/core/runtime.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
"""Runtime helpers for simulated pipeline steps."""

from __future__ import annotations

import asyncio
from typing import Awaitable, Callable


async def simulate_runtime(seconds: float = 1.0) -> None:
"""Sleep in small increments to emulate work without blocking tests."""

# Sleep using small intervals to keep control responsive while keeping
# the total runtime deterministic for tests.
remaining = float(seconds)
interval = 0.1
while remaining > 0:
await asyncio.sleep(min(interval, remaining))
remaining -= interval


def bind_runtime(seconds: float) -> Callable[[], Awaitable[None]]:
"""Return a coroutine factory that simulates work for ``seconds`` seconds."""

async def _runner() -> None:
await simulate_runtime(seconds)

return _runner
Loading