Skip to content

Ingestion Services

Generic Ingestion

pluginlake.api.services.ingestion

Ingestion service — validates, stores, and triggers processing of uploaded files.

Orchestrates the flow from file upload to raw storage layer and Dagster job triggering, with proper validation and error handling.

IngestionError

Bases: Exception

Raised when file ingestion fails due to validation or storage errors.

Source code in src/pluginlake/api/services/ingestion.py
class IngestionError(Exception):
    """Raised when file ingestion fails due to validation or storage errors."""

IngestionResult dataclass

Outcome of a file ingestion operation.

Attributes:

Name Type Description
file_id str

Unique identifier assigned to the uploaded file.

filename str

Original filename of the upload.

dataset str

Target dataset name.

file_path str

Absolute path where the file was stored.

size_bytes int

Size of the stored file in bytes.

dagster_run_id str | None

Dagster run ID if a job was triggered, else None.

status str

Overall status of the ingestion.

message str

Human-readable status message.

Source code in src/pluginlake/api/services/ingestion.py
@dataclass
class IngestionResult:
    """Outcome of a file ingestion operation.

    Attributes:
        file_id: Unique identifier assigned to the uploaded file.
        filename: Original filename of the upload.
        dataset: Target dataset name.
        file_path: Absolute path where the file was stored.
        size_bytes: Size of the stored file in bytes.
        dagster_run_id: Dagster run ID if a job was triggered, else None.
        status: Overall status of the ingestion.
        message: Human-readable status message.
    """

    file_id: str
    filename: str
    dataset: str
    file_path: str
    size_bytes: int
    dagster_run_id: str | None
    status: str
    message: str

IngestionService

Handles file upload validation, storage, and Dagster triggering.

Parameters:

Name Type Description Default
storage_manager StorageLayerManager | None

Manages the medallion storage layers.

None
dagster_client DagsterClient | None

Client for triggering Dagster jobs.

None
settings IngestionSettings | None

Ingestion-specific configuration.

None
Source code in src/pluginlake/api/services/ingestion.py
class IngestionService:
    """Handles file upload validation, storage, and Dagster triggering.

    Args:
        storage_manager: Manages the medallion storage layers.
        dagster_client: Client for triggering Dagster jobs.
        settings: Ingestion-specific configuration.
    """

    def __init__(
        self,
        storage_manager: StorageLayerManager | None = None,
        dagster_client: DagsterClient | None = None,
        settings: IngestionSettings | None = None,
        ingestion_log_dir: Path | None = None,
    ) -> None:
        """Initialise the service with storage, Dagster, and settings."""
        self._storage = storage_manager
        self._dagster = dagster_client
        self._settings = settings or IngestionSettings()
        self._ingestion_log_dir = ingestion_log_dir

    async def ingest_file(self, file: UploadFile, dataset: str) -> IngestionResult:
        """Validate, store, and trigger processing for an uploaded file.

        Args:
            file: The uploaded file from the HTTP request.
            dataset: Target dataset name for organising the file.

        Returns:
            An IngestionResult describing the outcome.

        Raises:
            IngestionError: If validation fails (extension, size).
        """
        filename = file.filename or "unnamed"
        file_id = uuid.uuid4().hex[:12]

        logger.info("Ingesting file %r for dataset %r (file_id=%s)", filename, dataset, file_id)

        self._validate_extension(filename)
        content = await self._read_and_validate_size(file, filename)
        file_path = self._store_file(content, filename, dataset, file_id)
        dagster_run_id = await self._trigger_dagster(file_path, dataset, filename)

        status = "completed" if dagster_run_id else "stored"
        message = (
            f"File ingested and Dagster job triggered (run_id={dagster_run_id})."
            if dagster_run_id
            else "File ingested successfully. Dagster job trigger was skipped or failed."
        )

        logger.info(
            "Ingestion %s: file_id=%s, dataset=%s, path=%s, dagster_run=%s",
            status,
            file_id,
            dataset,
            file_path,
            dagster_run_id,
        )

        result = IngestionResult(
            file_id=file_id,
            filename=filename,
            dataset=dataset,
            file_path=str(file_path),
            size_bytes=len(content),
            dagster_run_id=dagster_run_id,
            status=status,
            message=message,
        )

        if self._ingestion_log_dir is not None:
            write_ingestion_log(
                self._ingestion_log_dir,
                file_id=result.file_id,
                filename=result.filename,
                dataset=result.dataset,
                size_bytes=result.size_bytes,
                status=result.status,
                dagster_run_id=result.dagster_run_id,
            )

        return result

    def _validate_extension(self, filename: str) -> None:
        """Check that the file extension is in the allowed list."""
        suffix = Path(filename).suffix.lower()
        if suffix not in self._settings.allowed_extensions:
            msg = f"File extension {suffix!r} is not allowed. Allowed: {', '.join(self._settings.allowed_extensions)}"
            raise IngestionError(msg)

    async def _read_and_validate_size(self, file: UploadFile, filename: str) -> bytes:
        """Read the full file content and validate against the size limit."""
        content = await file.read()
        if len(content) > self._settings.max_file_size_bytes:
            msg = (
                f"File {filename!r} ({len(content)} bytes) exceeds the maximum "
                f"allowed size of {self._settings.max_file_size_mb} MB."
            )
            raise IngestionError(msg)
        return content

    def _store_file(
        self,
        content: bytes,
        filename: str,
        dataset: str,
        file_id: str,
    ) -> Path:
        """Write the file content to the raw storage layer.

        Files are stored as: ``raw/{dataset}/{timestamp}_{file_id}_{filename}``
        """
        if self._storage is None:
            msg = "_store_file requires storage_manager"
            raise TypeError(msg)
        timestamp = datetime.now(tz=UTC).strftime("%Y%m%dT%H%M%S")
        stored_name = f"{timestamp}_{file_id}_{filename}"
        dest = self._storage.get_path(StorageLayer.RAW, dataset, stored_name)

        dest.parent.mkdir(parents=True, exist_ok=True)
        dest.write_bytes(content)

        logger.info("Stored file at %s (%d bytes)", dest, len(content))
        return dest

    async def _trigger_dagster(
        self,
        file_path: Path,
        dataset: str,
        filename: str,
    ) -> str | None:
        """Attempt to trigger a Dagster job for the ingested file.

        Returns the run ID on success, or None if triggering fails.
        Failures are logged but do not raise — the file is already stored.
        """
        if self._dagster is None:
            msg = "_trigger_dagster requires dagster_client"
            raise TypeError(msg)
        try:
            result = await self._dagster.trigger_job(
                job_name=self._settings.dagster_job_name,
                run_config={
                    "ops": {
                        "process_raw_file": {
                            "config": {
                                "file_path": str(file_path),
                                "dataset": dataset,
                                "filename": filename,
                            },
                        },
                    },
                },
            )
        except DagsterClientError:
            logger.warning(
                "Failed to trigger Dagster job for %r (dataset=%r). "
                "File was stored successfully — processing can be triggered manually.",
                filename,
                dataset,
                exc_info=True,
            )
            return None
        else:
            return result.run_id

__init__(storage_manager=None, dagster_client=None, settings=None, ingestion_log_dir=None)

Initialise the service with storage, Dagster, and settings.

Source code in src/pluginlake/api/services/ingestion.py
def __init__(
    self,
    storage_manager: StorageLayerManager | None = None,
    dagster_client: DagsterClient | None = None,
    settings: IngestionSettings | None = None,
    ingestion_log_dir: Path | None = None,
) -> None:
    """Initialise the service with storage, Dagster, and settings."""
    self._storage = storage_manager
    self._dagster = dagster_client
    self._settings = settings or IngestionSettings()
    self._ingestion_log_dir = ingestion_log_dir

ingest_file(file, dataset) async

Validate, store, and trigger processing for an uploaded file.

Parameters:

Name Type Description Default
file UploadFile

The uploaded file from the HTTP request.

required
dataset str

Target dataset name for organising the file.

required

Returns:

Type Description
IngestionResult

An IngestionResult describing the outcome.

Raises:

Type Description
IngestionError

If validation fails (extension, size).

Source code in src/pluginlake/api/services/ingestion.py
async def ingest_file(self, file: UploadFile, dataset: str) -> IngestionResult:
    """Validate, store, and trigger processing for an uploaded file.

    Args:
        file: The uploaded file from the HTTP request.
        dataset: Target dataset name for organising the file.

    Returns:
        An IngestionResult describing the outcome.

    Raises:
        IngestionError: If validation fails (extension, size).
    """
    filename = file.filename or "unnamed"
    file_id = uuid.uuid4().hex[:12]

    logger.info("Ingesting file %r for dataset %r (file_id=%s)", filename, dataset, file_id)

    self._validate_extension(filename)
    content = await self._read_and_validate_size(file, filename)
    file_path = self._store_file(content, filename, dataset, file_id)
    dagster_run_id = await self._trigger_dagster(file_path, dataset, filename)

    status = "completed" if dagster_run_id else "stored"
    message = (
        f"File ingested and Dagster job triggered (run_id={dagster_run_id})."
        if dagster_run_id
        else "File ingested successfully. Dagster job trigger was skipped or failed."
    )

    logger.info(
        "Ingestion %s: file_id=%s, dataset=%s, path=%s, dagster_run=%s",
        status,
        file_id,
        dataset,
        file_path,
        dagster_run_id,
    )

    result = IngestionResult(
        file_id=file_id,
        filename=filename,
        dataset=dataset,
        file_path=str(file_path),
        size_bytes=len(content),
        dagster_run_id=dagster_run_id,
        status=status,
        message=message,
    )

    if self._ingestion_log_dir is not None:
        write_ingestion_log(
            self._ingestion_log_dir,
            file_id=result.file_id,
            filename=result.filename,
            dataset=result.dataset,
            size_bytes=result.size_bytes,
            status=result.status,
            dagster_run_id=result.dagster_run_id,
        )

    return result

OMOP CSV Ingestion

pluginlake.api.services.omop_ingestion

OMOP ingestion service — CSV upload path.

OmopCsvIngestionService

Bases: IngestionService

Ingestion service for OMOP CSV uploads.

Stores the uploaded CSV to OMOP_RAW_DATA_DIR/{dataset}.csv and triggers omop_ingest_job for the selected asset via Dagster.

Parameters:

Name Type Description Default
dagster_client DagsterClient

Client for triggering Dagster jobs.

required
omop_settings OMOPSettings

OMOP-specific configuration (provides raw_data_dir).

required
settings IngestionSettings | None

Ingestion-specific configuration (file size limits, etc.).

None
Source code in src/pluginlake/api/services/omop_ingestion.py
class OmopCsvIngestionService(IngestionService):
    """Ingestion service for OMOP CSV uploads.

    Stores the uploaded CSV to ``OMOP_RAW_DATA_DIR/{dataset}.csv`` and
    triggers ``omop_ingest_job`` for the selected asset via Dagster.

    Args:
        dagster_client: Client for triggering Dagster jobs.
        omop_settings: OMOP-specific configuration (provides raw_data_dir).
        settings: Ingestion-specific configuration (file size limits, etc.).
    """

    def __init__(
        self,
        dagster_client: DagsterClient,
        omop_settings: OMOPSettings,
        settings: IngestionSettings | None = None,
        ingestion_log_dir: Path | None = None,
    ) -> None:
        """Initialise with a Dagster client, OMOP settings, and optional ingestion settings."""
        super().__init__(dagster_client=dagster_client, settings=settings, ingestion_log_dir=ingestion_log_dir)
        self._omop_settings = omop_settings

    def _validate_extension(self, filename: str) -> None:
        if Path(filename).suffix.lower() != ".csv":
            msg = "OMOP table uploads must be .csv files."
            raise IngestionError(msg)

    def _store_file(self, content: bytes, filename: str, dataset: str, file_id: str) -> Path:  # noqa: ARG002 — signature required by IngestionService; OMOP stores by dataset name only
        dest = self._omop_settings.raw_data_dir / f"{dataset}.csv"
        dest.parent.mkdir(parents=True, exist_ok=True)
        dest.write_bytes(content)
        return dest

    async def _trigger_dagster(self, file_path: Path, dataset: str, filename: str) -> str | None:  # noqa: ARG002 — signature required by IngestionService; OMOP triggers by dataset name only
        if self._dagster is None:
            msg = "_trigger_dagster requires dagster_client"
            raise TypeError(msg)
        try:
            result = await self._dagster.trigger_job(
                job_name="omop_ingest_job",
                asset_selection=[["omop", dataset]],
            )
        except DagsterClientError as exc:
            logger.warning("Failed to trigger Dagster for omop/%s: %s", dataset, exc)
            return None
        else:
            return result.run_id

__init__(dagster_client, omop_settings, settings=None, ingestion_log_dir=None)

Initialise with a Dagster client, OMOP settings, and optional ingestion settings.

Source code in src/pluginlake/api/services/omop_ingestion.py
def __init__(
    self,
    dagster_client: DagsterClient,
    omop_settings: OMOPSettings,
    settings: IngestionSettings | None = None,
    ingestion_log_dir: Path | None = None,
) -> None:
    """Initialise with a Dagster client, OMOP settings, and optional ingestion settings."""
    super().__init__(dagster_client=dagster_client, settings=settings, ingestion_log_dir=ingestion_log_dir)
    self._omop_settings = omop_settings

FHIR NDJSON Ingestion

pluginlake.api.services.fhir_ingestion

FHIR ingestion service — NDJSON upload path.

FhirNdjsonIngestionService

Bases: IngestionService

Ingestion service for FHIR NDJSON uploads.

Appends the uploaded NDJSON to FHIR_RAW_DATA_DIR/{dataset}.ndjson and triggers fhir_ingest_job for the selected asset via Dagster.

Parameters:

Name Type Description Default
dagster_client DagsterClient

Client for triggering Dagster jobs.

required
fhir_settings FHIRSettings

FHIR-specific configuration (provides raw_data_dir).

required
settings IngestionSettings | None

Ingestion-specific configuration (file size limits, etc.).

None
Source code in src/pluginlake/api/services/fhir_ingestion.py
class FhirNdjsonIngestionService(IngestionService):
    """Ingestion service for FHIR NDJSON uploads.

    Appends the uploaded NDJSON to ``FHIR_RAW_DATA_DIR/{dataset}.ndjson`` and
    triggers ``fhir_ingest_job`` for the selected asset via Dagster.

    Args:
        dagster_client: Client for triggering Dagster jobs.
        fhir_settings: FHIR-specific configuration (provides raw_data_dir).
        settings: Ingestion-specific configuration (file size limits, etc.).
    """

    def __init__(
        self,
        dagster_client: DagsterClient,
        fhir_settings: FHIRSettings,
        settings: IngestionSettings | None = None,
        ingestion_log_dir: Path | None = None,
    ) -> None:
        """Initialise with a Dagster client, FHIR settings, and optional ingestion settings."""
        super().__init__(dagster_client=dagster_client, settings=settings, ingestion_log_dir=ingestion_log_dir)
        self._fhir_settings = fhir_settings

    def _validate_extension(self, filename: str) -> None:
        if Path(filename).suffix.lower() != ".ndjson":
            msg = "FHIR uploads must be .ndjson files."
            raise IngestionError(msg)

    def _store_file(self, content: bytes, filename: str, dataset: str, file_id: str) -> Path:  # noqa: ARG002 — signature required by IngestionService; FHIR stores by dataset name only
        dest = self._fhir_settings.raw_data_dir / f"{dataset}.ndjson"
        dest.parent.mkdir(parents=True, exist_ok=True)
        needs_separator = dest.exists() and dest.stat().st_size > 0
        with dest.open("ab") as f:
            if needs_separator:
                f.write(b"\n")
            f.write(content)
        return dest

    async def _trigger_dagster(self, file_path: Path, dataset: str, filename: str) -> str | None:  # noqa: ARG002 — signature required by IngestionService; FHIR triggers by dataset name only
        if self._dagster is None:
            msg = "_trigger_dagster requires dagster_client"
            raise TypeError(msg)
        try:
            asset_selection: list[list[str]] = [["fhir_raw", dataset]]
            omop_table = FHIR_TO_OMOP_TABLE.get(dataset)
            if omop_table:
                asset_selection.append(["fhir_omop_raw", omop_table])
                asset_selection.append(["omop", omop_table])
            result = await self._dagster.trigger_job(
                job_name="fhir_ingest_job",
                asset_selection=asset_selection,
            )
        except DagsterClientError as exc:
            logger.warning("Failed to trigger Dagster for fhir_raw/%s: %s", dataset, exc)
            return None
        else:
            return result.run_id

__init__(dagster_client, fhir_settings, settings=None, ingestion_log_dir=None)

Initialise with a Dagster client, FHIR settings, and optional ingestion settings.

Source code in src/pluginlake/api/services/fhir_ingestion.py
def __init__(
    self,
    dagster_client: DagsterClient,
    fhir_settings: FHIRSettings,
    settings: IngestionSettings | None = None,
    ingestion_log_dir: Path | None = None,
) -> None:
    """Initialise with a Dagster client, FHIR settings, and optional ingestion settings."""
    super().__init__(dagster_client=dagster_client, settings=settings, ingestion_log_dir=ingestion_log_dir)
    self._fhir_settings = fhir_settings