Skip to content

modosaic.core.pipeline

modosaic.core.pipeline

Pipeline

Pipeline(dataset, modalities, experiment=None)

Sequential multimodal generation and validation pipeline.

The pipeline iterates over input image records, runs each configured modality in order, applies validators and optional constraints, and saves artifacts only for modalities whose constraints pass.

Attributes:

Name Type Description
dataset Iterable[ImageRecord]

Iterable source of image records.

modalities list[Modality[Any]]

Ordered modality definitions to run for each sample.

experiment ExperimentService

Service used to persist accepted artifacts.

Initialize a pipeline.

Parameters:

Name Type Description Default
dataset Iterable[ImageRecord]

Image records to process.

required
modalities Iterable[Modality[Any]]

Ordered modality definitions. Later modalities can read accepted outputs from earlier modalities through validator dependencies.

required
experiment ExperimentService | None

Optional experiment service. A default service writing to experiments/<timestamp> is created when omitted.

None
Source code in modosaic/core/pipeline.py
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
def __init__(
        self,
        dataset: Iterable[ImageRecord],
        modalities: Iterable[Modality[Any]],
        experiment: ExperimentService | None = None,
) -> None:
    """Initialize a pipeline.

    Args:
        dataset: Image records to process.
        modalities: Ordered modality definitions. Later modalities can read
            accepted outputs from earlier modalities through validator
            dependencies.
        experiment: Optional experiment service. A default service writing
            to `experiments/<timestamp>` is created when omitted.
    """
    self.dataset = dataset
    self.modalities = list(modalities)
    self.experiment = experiment or ExperimentService()

run

run(limit=None)

Run the pipeline over the configured dataset.

Parameters:

Name Type Description Default
limit int | None

Optional maximum number of samples to process.

None

Returns:

Type Description
list[PipelineSampleResult]

Results for every processed sample.

Source code in modosaic/core/pipeline.py
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
def run(self, limit: int | None = None) -> list[PipelineSampleResult]:
    """Run the pipeline over the configured dataset.

    Args:
        limit: Optional maximum number of samples to process.

    Returns:
        Results for every processed sample.
    """
    results: list[PipelineSampleResult] = []

    for idx, record in enumerate(self.dataset):
        if limit is not None and idx >= limit:
            logger.info(f"Limit reached, stopping pipeline")
            break

        logger.debug(f"Starting pipeline for record sample {record.sample_id}")
        results.append(self.run_sample(record))

    return results

run_sample

run_sample(record)

Run all configured modalities for one sample.

Parameters:

Name Type Description Default
record ImageRecord

Image record to process.

required

Returns:

Type Description
PipelineSampleResult

Generated outputs, validation results, and saved artifact paths for

PipelineSampleResult

the sample.

Source code in modosaic/core/pipeline.py
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
def run_sample(self, record: ImageRecord) -> PipelineSampleResult:
    """Run all configured modalities for one sample.

    Args:
        record: Image record to process.

    Returns:
        Generated outputs, validation results, and saved artifact paths for
        the sample.
    """
    generated_modalities: dict[Modalities, Any] = {}
    validations: dict[Modalities, list[ValidationResult[Any]]] = {}
    artifact_paths: list[Path] = []

    for modality in self.modalities:
        logger.debug(f"Processing modality: {modality.modality} for record sample {record.sample_id}")

        generated = modality.generate(record)

        modality_validations = modality.validate(
            record=record,
            generated=generated,
            generated_modalities=generated_modalities,
        )
        validations[modality.modality] = modality_validations

        if self._has_failed_constraint(modality_validations):
            logger.info(
                f"Discarding modality {modality.modality} for record sample {record.sample_id} "
                f"after failed validation"
            )
            continue

        generated_modalities[modality.modality] = generated

        artifacts = modality.postprocess(record, generated)
        logger.debug(f"Generated {len(artifacts)} artifacts for modality {modality.modality} on record sample {record.sample_id}")
        artifact_paths.extend(self.experiment.save_artifacts(artifacts))
        logger.debug(
            f"Saved {len(artifacts)} artifacts for modality {modality.modality} "
            f"on record sample {record.sample_id}"
        )

        if modality_validations:
            validation_artifact = self._validation_artifact(
                record=record,
                modality=modality.modality,
                validations=modality_validations,
            )
            artifact_paths.append(self.experiment.save_artifact(validation_artifact))

    if self.modalities and not generated_modalities:
        logger.info(f"Discarded record sample {record.sample_id}; no modalities passed validation")
    else:
        logger.info(
            f"Completed record sample {record.sample_id} with {len(generated_modalities)} modalities "
            f"and {len(artifact_paths)} artifacts"
        )

    return PipelineSampleResult(
        record=record,
        generated=generated_modalities,
        validations=validations,
        artifact_paths=artifact_paths,
    )

PipelineSampleResult dataclass

PipelineSampleResult(record, generated, validations, artifact_paths)

Outputs collected after processing one image record.

Attributes:

Name Type Description
record ImageRecord

Original image record processed by the pipeline.

generated Mapping[Modalities, Any]

Accepted generated outputs keyed by modality.

validations Mapping[Modalities, list[ValidationResult[Any]]]

Validation results keyed by modality.

artifact_paths list[Path]

Files written for accepted modality outputs and validation reports.