Skip to content

Pipeline

vartriage.Pipeline

Top-level orchestrator for the variant prioritization pipeline.

Connects all stages and manages configuration validation, streaming variants through the pipeline while maintaining memory bounds suitable for whole-genome scale datasets (4M+ variants, <2GB RSS).

The pipeline validates all configuration parameters at construction time (fail-fast). The run method then executes the full processing chain: parsing → filtering → annotation → prioritization → classification → report generation.

Parameters

config : PipelineConfig Complete configuration for all pipeline stages.

Raises

ValueError If any configuration parameter is invalid (out-of-range thresholds, invalid batch sizes, unsupported formats). FileNotFoundError If any required reference file is missing.

Source code in vartriage/pipeline.py
  38
  39
  40
  41
  42
  43
  44
  45
  46
  47
  48
  49
  50
  51
  52
  53
  54
  55
  56
  57
  58
  59
  60
  61
  62
  63
  64
  65
  66
  67
  68
  69
  70
  71
  72
  73
  74
  75
  76
  77
  78
  79
  80
  81
  82
  83
  84
  85
  86
  87
  88
  89
  90
  91
  92
  93
  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
 159
 160
 161
 162
 163
 164
 165
 166
 167
 168
 169
 170
 171
 172
 173
 174
 175
 176
 177
 178
 179
 180
 181
 182
 183
 184
 185
 186
 187
 188
 189
 190
 191
 192
 193
 194
 195
 196
 197
 198
 199
 200
 201
 202
 203
 204
 205
 206
 207
 208
 209
 210
 211
 212
 213
 214
 215
 216
 217
 218
 219
 220
 221
 222
 223
 224
 225
 226
 227
 228
 229
 230
 231
 232
 233
 234
 235
 236
 237
 238
 239
 240
 241
 242
 243
 244
 245
 246
 247
 248
 249
 250
 251
 252
 253
 254
 255
 256
 257
 258
 259
 260
 261
 262
 263
 264
 265
 266
 267
 268
 269
 270
 271
 272
 273
 274
 275
 276
 277
 278
 279
 280
 281
 282
 283
 284
 285
 286
 287
 288
 289
 290
 291
 292
 293
 294
 295
 296
 297
 298
 299
 300
 301
 302
 303
 304
 305
 306
 307
 308
 309
 310
 311
 312
 313
 314
 315
 316
 317
 318
 319
 320
 321
 322
 323
 324
 325
 326
 327
 328
 329
 330
 331
 332
 333
 334
 335
 336
 337
 338
 339
 340
 341
 342
 343
 344
 345
 346
 347
 348
 349
 350
 351
 352
 353
 354
 355
 356
 357
 358
 359
 360
 361
 362
 363
 364
 365
 366
 367
 368
 369
 370
 371
 372
 373
 374
 375
 376
 377
 378
 379
 380
 381
 382
 383
 384
 385
 386
 387
 388
 389
 390
 391
 392
 393
 394
 395
 396
 397
 398
 399
 400
 401
 402
 403
 404
 405
 406
 407
 408
 409
 410
 411
 412
 413
 414
 415
 416
 417
 418
 419
 420
 421
 422
 423
 424
 425
 426
 427
 428
 429
 430
 431
 432
 433
 434
 435
 436
 437
 438
 439
 440
 441
 442
 443
 444
 445
 446
 447
 448
 449
 450
 451
 452
 453
 454
 455
 456
 457
 458
 459
 460
 461
 462
 463
 464
 465
 466
 467
 468
 469
 470
 471
 472
 473
 474
 475
 476
 477
 478
 479
 480
 481
 482
 483
 484
 485
 486
 487
 488
 489
 490
 491
 492
 493
 494
 495
 496
 497
 498
 499
 500
 501
 502
 503
 504
 505
 506
 507
 508
 509
 510
 511
 512
 513
 514
 515
 516
 517
 518
 519
 520
 521
 522
 523
 524
 525
 526
 527
 528
 529
 530
 531
 532
 533
 534
 535
 536
 537
 538
 539
 540
 541
 542
 543
 544
 545
 546
 547
 548
 549
 550
 551
 552
 553
 554
 555
 556
 557
 558
 559
 560
 561
 562
 563
 564
 565
 566
 567
 568
 569
 570
 571
 572
 573
 574
 575
 576
 577
 578
 579
 580
 581
 582
 583
 584
 585
 586
 587
 588
 589
 590
 591
 592
 593
 594
 595
 596
 597
 598
 599
 600
 601
 602
 603
 604
 605
 606
 607
 608
 609
 610
 611
 612
 613
 614
 615
 616
 617
 618
 619
 620
 621
 622
 623
 624
 625
 626
 627
 628
 629
 630
 631
 632
 633
 634
 635
 636
 637
 638
 639
 640
 641
 642
 643
 644
 645
 646
 647
 648
 649
 650
 651
 652
 653
 654
 655
 656
 657
 658
 659
 660
 661
 662
 663
 664
 665
 666
 667
 668
 669
 670
 671
 672
 673
 674
 675
 676
 677
 678
 679
 680
 681
 682
 683
 684
 685
 686
 687
 688
 689
 690
 691
 692
 693
 694
 695
 696
 697
 698
 699
 700
 701
 702
 703
 704
 705
 706
 707
 708
 709
 710
 711
 712
 713
 714
 715
 716
 717
 718
 719
 720
 721
 722
 723
 724
 725
 726
 727
 728
 729
 730
 731
 732
 733
 734
 735
 736
 737
 738
 739
 740
 741
 742
 743
 744
 745
 746
 747
 748
 749
 750
 751
 752
 753
 754
 755
 756
 757
 758
 759
 760
 761
 762
 763
 764
 765
 766
 767
 768
 769
 770
 771
 772
 773
 774
 775
 776
 777
 778
 779
 780
 781
 782
 783
 784
 785
 786
 787
 788
 789
 790
 791
 792
 793
 794
 795
 796
 797
 798
 799
 800
 801
 802
 803
 804
 805
 806
 807
 808
 809
 810
 811
 812
 813
 814
 815
 816
 817
 818
 819
 820
 821
 822
 823
 824
 825
 826
 827
 828
 829
 830
 831
 832
 833
 834
 835
 836
 837
 838
 839
 840
 841
 842
 843
 844
 845
 846
 847
 848
 849
 850
 851
 852
 853
 854
 855
 856
 857
 858
 859
 860
 861
 862
 863
 864
 865
 866
 867
 868
 869
 870
 871
 872
 873
 874
 875
 876
 877
 878
 879
 880
 881
 882
 883
 884
 885
 886
 887
 888
 889
 890
 891
 892
 893
 894
 895
 896
 897
 898
 899
 900
 901
 902
 903
 904
 905
 906
 907
 908
 909
 910
 911
 912
 913
 914
 915
 916
 917
 918
 919
 920
 921
 922
 923
 924
 925
 926
 927
 928
 929
 930
 931
 932
 933
 934
 935
 936
 937
 938
 939
 940
 941
 942
 943
 944
 945
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
class Pipeline:
    """Top-level orchestrator for the variant prioritization pipeline.

    Connects all stages and manages configuration validation, streaming
    variants through the pipeline while maintaining memory bounds suitable
    for whole-genome scale datasets (4M+ variants, <2GB RSS).

    The pipeline validates all configuration parameters at construction
    time (fail-fast). The ``run`` method then executes the full processing
    chain: parsing → filtering → annotation → prioritization →
    classification → report generation.

    Parameters
    ----------
    config : PipelineConfig
        Complete configuration for all pipeline stages.

    Raises
    ------
    ValueError
        If any configuration parameter is invalid (out-of-range thresholds,
        invalid batch sizes, unsupported formats).
    FileNotFoundError
        If any required reference file is missing.
    """

    def __init__(self, config: PipelineConfig) -> None:
        self._config = config
        self._warning_accumulator = WarningAccumulator(config.missing_data)
        self._validate_config(config)

        # NMD-escape lookup is built lazily from the GTF and cached. The
        # sentinel False means "not yet built"; None means "built, unavailable".
        self._nmd_lookup_cached: object | bool | None = False

        # Gene-disease linkage annotator: constructed once, reused across runs
        self._gene_knowledge_annotator = None  # type: GeneKnowledgeAnnotator | None
        if config.knowledge is not None:
            from vartriage.knowledge.annotator import (
                GeneKnowledgeAnnotator,
            )

            self._gene_knowledge_annotator = GeneKnowledgeAnnotator(config.knowledge)

    @property
    def _mito_enabled(self) -> bool:
        """Whether mitochondrial analysis is active."""
        if self._config.mito is not None:
            return self._config.mito.enabled
        # Auto-enabled by default (no explicit config needed)
        return True

    @property
    def mito_results(self) -> list[Any] | None:
        """Access mitochondrial classification results from the last run.

        Returns None if mitochondrial analysis was disabled or no chrM
        variants were present in the input.
        """
        return getattr(self, "_mito_results", None)

    @property
    def warning_accumulator(self) -> WarningAccumulator:
        """Access the warning accumulator tracking MissingDataWarnings.

        Returns
        -------
        WarningAccumulator
            The shared warning accumulator for the current pipeline run.
        """
        return self._warning_accumulator

    def _build_stages(
        self,
    ) -> tuple[
        QualityFilter,
        AnnotationEngine | None,
        object,
        PrioritizationEngine,
        ACMGClassifier,
        ReportGenerator,
    ]:
        """Construct all pipeline stage objects."""
        quality_filter = QualityFilter(self._config.quality_filter)

        annotation_engine: AnnotationEngine | None = None
        api_annotation_engine = None

        if self._config.api is not None:
            api_annotation_engine = self._build_api_annotation_engine()

        if api_annotation_engine is None and self._config.annotation is not None:
            annotation_engine = AnnotationEngine(self._config.annotation)

        # Inject remote gnomAD backend when no local gnomAD and remote is configured
        self._inject_remote_gnomad(annotation_engine)

        prioritization_engine = PrioritizationEngine(
            self._config.prioritization, remote_config=self._config.remote
        )
        acmg_classifier = self._build_acmg_classifier()
        report_generator = ReportGenerator(
            self._config.report,
            clinical_config=self._config.clinical_report,
            reference_checksums=(
                self._compute_reference_checksums()
                if self._config.clinical_report is not None
                else None
            ),
        )
        return (
            quality_filter,
            annotation_engine,
            api_annotation_engine,
            prioritization_engine,
            acmg_classifier,
            report_generator,
        )

    def _run_qc_pass(self, vcf_path: Path) -> Any:
        """Run the QC pre-flight pass if configured.

        Reads the VCF independently (separate file handle) to compute
        population-level summary statistics. When --strict-qc is active,
        raises SystemExit(3) on any FAIL-level metric.

        Returns the QCReport or None if QC is skipped.
        """
        qc_config = self._config.qc
        if qc_config is None or qc_config.skip:
            self._qc_report = None
            return None

        from vartriage.qc.metrics import compute_qc_metrics
        from vartriage.qc.report import print_qc_stderr
        from vartriage.qc.validator import QCStatus, QCValidator

        sample_id = qc_config.sample_id
        if sample_id is None and self._config.sample is not None:
            sample_id = self._config.sample.sample_name

        logger.info("Running QC pre-flight pass (assay: %s)", qc_config.assay_type)
        metrics = compute_qc_metrics(vcf_path, sample_id=sample_id)

        validator = QCValidator(qc_config)
        report = validator.validate(metrics)

        print_qc_stderr(report)
        self._qc_report = report

        if qc_config.strict and report.overall_status == QCStatus.FAIL:
            failed_checks = [c for c in report.checks if c.status == QCStatus.FAIL]
            fail_msg = "; ".join(c.message for c in failed_checks)
            logger.error("QC FAIL (strict mode): %s", fail_msg)
            import sys

            print(
                f"Error: QC pre-flight FAILED (strict mode): {fail_msg}",
                file=sys.stderr,
            )
            sys.exit(3)

        return report

    @property
    def qc_report(self) -> Any:
        """Access the QC report from the last run. None if QC was skipped."""
        return getattr(self, "_qc_report", None)

    def _generate_report(
        self,
        report_generator: ReportGenerator,
        classified: Iterator[ClassifiedVariant],
        output_path: Path,
        vcf_path: Path,
        mito_results: list[Any] | None = None,
    ) -> Path:
        """Dispatch report generation for VCF and non-VCF formats.

        When mito_results are present, uses mito-aware writers routed
        through ReportGenerator's atomic write logic (temp file + rename).
        """
        self._mito_results = mito_results
        fmt = self._config.report.output_format

        if fmt == "vcf":
            return report_generator.generate(classified, output_path, vcf_path)

        # For JSON/CSV with mito results, use mito-aware writers via
        # the atomic write strategy (temp file + os.replace)
        if mito_results and fmt in ("json", "csv"):
            return self._generate_mito_report(
                classified, output_path, fmt, mito_results
            )

        # Clinical formats always carry the QC report (when QC ran) so the
        # Sample Quality Control section renders regardless of mito presence.
        if fmt.startswith("clinical-"):
            return report_generator.generate(
                classified,
                output_path,
                mito_results=mito_results,
                qc_report=self.qc_report,
            )

        return report_generator.generate(classified, output_path)

    def _generate_mito_report(
        self,
        classified: Iterator[ClassifiedVariant],
        output_path: Path,
        fmt: str,
        mito_results: list[Any],
    ) -> Path:
        """Write mito-aware JSON/CSV with atomic temp-file strategy."""
        import os
        import tempfile

        from vartriage._internal.path_safety import resolve_path as _rp

        resolved = _rp(output_path)
        resolved.parent.mkdir(parents=True, exist_ok=True)

        tmp_fd, tmp_name = tempfile.mkstemp(
            dir=resolved.parent,
            prefix=".report_mito_",
            suffix=f".{fmt}.tmp",
        )
        os.close(tmp_fd)
        tmp_path = Path(tmp_name)

        try:
            if fmt == "json":
                from vartriage.reporting.json_writer import write_json_with_mito

                write_json_with_mito(classified, tmp_path, mito_results)
            else:
                from vartriage.reporting.csv_writer import write_csv_with_mito

                write_csv_with_mito(classified, tmp_path, mito_results)

            os.replace(str(tmp_path), str(resolved))
        except Exception:
            tmp_path.unlink(missing_ok=True)
            raise

        return resolved

    def run(
        self, vcf_path: Path | None = None, output_path: Path | None = None
    ) -> Path:
        """Execute the full variant prioritization pipeline.

        Wires stages sequentially: VCFParser → QualityFilter →
        AnnotationEngine → PrioritizationEngine → ACMGClassifier →
        ReportGenerator. Processes data in batches to maintain peak RSS
        below 2GB for whole-genome scale files (4M+ variants).

        Parameters
        ----------
        vcf_path : Path, optional
            Path to the input VCF file (.vcf or .vcf.gz). If None, uses
            the path from config.
        output_path : Path, optional
            Path for the output report. If None, uses the path from config.

        Returns
        -------
        Path
            Path to the generated output report file.

        Raises
        ------
        FileNotFoundError
            If the VCF file or any reference file is missing.
        ParseError
            If the VCF file has malformed headers or data lines.
        IOError
            If report generation fails due to a write error.
        """
        effective_vcf_path = vcf_path or self._config.vcf_path
        effective_output_path = output_path or self._config.output_path

        self._warning_accumulator.reset()

        logger.info(
            "Starting pipeline run: %s → %s", effective_vcf_path, effective_output_path
        )

        # QC pre-flight pass (reads VCF independently before annotation).
        # Result is stored on self._qc_report for report generation and
        # programmatic access via the qc_report property.
        self._run_qc_pass(effective_vcf_path)

        if self._config.clinical_report is not None:
            self._check_reference_checksums()

        (
            quality_filter,
            annotation_engine,
            api_annotation_engine,
            prioritization_engine,
            acmg_classifier,
            report_generator,
        ) = self._build_stages()

        # Per-sample AD/AF are needed for heteroplasmy and maternal
        # inheritance, so extract them whenever mito analysis is active, not
        # only when an inheritance or single-sample selection is configured.
        extract_samples = (
            self._config.inheritance is not None
            or self._config.sample is not None
            or self._mito_enabled
        )

        with VCFParser(
            effective_vcf_path,
            extract_samples=extract_samples,
        ) as parser:
            mito_variants, nuclear_iter = self._split_mito_variants(parser)

            annotated = self._build_annotated_stream(
                nuclear_iter,
                quality_filter,
                annotation_engine,
                api_annotation_engine,
                sample_names=parser.sample_names,
            )

            if self._config.gene_filter is not None:
                from vartriage.filter.gene_filter import GeneFilter

                annotated = GeneFilter(self._config.gene_filter).apply(annotated)

            if self._gene_knowledge_annotator is not None:
                annotated = self._gene_knowledge_annotator.annotate(annotated)

            scored = prioritization_engine.prioritize(annotated)

            if self._gene_knowledge_annotator is not None:
                scored = self._gene_knowledge_annotator.boost_scores(scored)

            classified = acmg_classifier.classify(scored)

            try:
                if self._mito_enabled:
                    classified_list = list(classified)
                    mito_results = self._run_mito_pipeline(mito_variants)
                    result_path = self._generate_report(
                        report_generator,
                        iter(classified_list),
                        effective_output_path,
                        effective_vcf_path,
                        mito_results=mito_results,
                    )
                else:
                    result_path = self._generate_report(
                        report_generator,
                        classified,
                        effective_output_path,
                        effective_vcf_path,
                        mito_results=None,
                    )

                if annotation_engine is not None:
                    self._warning_accumulator.add_batch(annotation_engine.warnings)
            finally:
                prioritization_engine.close()

        logger.info(
            "Pipeline completed. Missing data warnings: %d",
            self._warning_accumulator.total_count,
        )
        logger.info("Report written to: %s", result_path)
        return result_path

    def _split_mito_variants(
        self, parser: VCFParser
    ) -> tuple[list[Variant], Iterator[Variant]]:
        """Partition parser variants into mitochondrial and nuclear streams.

        Collects chrM variants into a list (small, typically <1000) while
        yielding nuclear variants lazily via a generator to avoid buffering
        millions of nuclear variants in memory.
        """
        if not self._mito_enabled:
            return [], iter(parser)

        mito: list[Variant] = []

        def _nuclear_generator() -> Iterator[Variant]:
            for variant in parser:
                if is_mitochondrial(variant.chrom):
                    mito.append(variant)
                else:
                    yield variant

        return mito, _nuclear_generator()

    def _run_mito_pipeline(self, mito_variants: list[Variant]) -> list[Any] | None:
        """Run the mitochondrial pipeline if variants are present.

        Caches the MitochondrialPipeline instance on first use so
        repeated calls (e.g. in cohort mode) don't reload TSV data.
        """
        if not (self._mito_enabled and mito_variants):
            return None
        from vartriage.mito.config import MitoConfig
        from vartriage.mito.pipeline import MitochondrialPipeline

        if not hasattr(self, "_mito_pipeline_instance"):
            mito_config = self._config.mito or MitoConfig()
            self._mito_pipeline_instance = MitochondrialPipeline(mito_config)

        results = self._mito_pipeline_instance.run(iter(mito_variants))
        logger.info("Mitochondrial pipeline: %d variants classified", len(results))
        return results

    def run_to_classification(
        self, vcf_path: Path | None = None
    ) -> Iterator[ClassifiedVariant]:
        """Execute the pipeline through classification without report generation.

        Runs VCF parsing, quality filtering, annotation, prioritization,
        and ACMG classification, yielding ClassifiedVariant objects. Skips
        report writing entirely. Used by CohortPipeline to collect per-sample
        results without generating throw-away files.

        Stage limitations vs full run():
            - SampleExtractor is not applied (cohort mode processes
              each VCF as a single-sample file already).
            - InheritanceFilter is not applied (trio analysis is a
              single-proband concern, not meaningful per-sample in a
              cohort where each VCF represents one individual).
            - SecondaryFindingsFilter is not applied (cohort mode
              aggregates by variant, not by clinical reporting context).

        These stages are intentionally excluded because cohort analysis
        operates on pre-separated per-sample VCFs. If your workflow
        requires sample extraction or inheritance filtering, run the
        full single-sample pipeline first and feed ClassifiedVariant
        results to CohortAggregator directly.

        Parameters
        ----------
        vcf_path : Path, optional
            Path to the input VCF file. If None, uses the path from config.

        Yields
        ------
        ClassifiedVariant
            Classified variants from the pipeline stages.
        """
        effective_vcf_path = vcf_path or self._config.vcf_path

        quality_filter = QualityFilter(self._config.quality_filter)

        annotation_engine: AnnotationEngine | None = None
        if self._config.annotation is not None:
            annotation_engine = AnnotationEngine(self._config.annotation)

        # Inject remote gnomAD for cohort mode
        self._inject_remote_gnomad(annotation_engine)

        prioritization_engine = PrioritizationEngine(
            self._config.prioritization, remote_config=self._config.remote
        )
        acmg_classifier = self._build_acmg_classifier()

        try:
            with VCFParser(effective_vcf_path) as parser:
                stream: Iterator[Variant] = iter(parser)

                if self._config.region_filter is not None:
                    from vartriage.filter.region_filter import RegionFilter

                    region_filter = RegionFilter(self._config.region_filter)
                    stream = region_filter.apply(stream)

                filtered = quality_filter.apply(stream)

                if annotation_engine is not None:
                    annotated = annotation_engine.annotate(filtered)
                else:
                    annotated = self._passthrough_annotation(filtered)

                if self._config.gene_filter is not None:
                    from vartriage.filter.gene_filter import GeneFilter

                    gene_filter = GeneFilter(self._config.gene_filter)
                    annotated = gene_filter.apply(annotated)

                # Gene-disease linkage in run_to_classification path
                if self._gene_knowledge_annotator is not None:
                    annotated = self._gene_knowledge_annotator.annotate(annotated)

                scored = prioritization_engine.prioritize(annotated)

                # Apply phenotype boost in classification path too
                if self._gene_knowledge_annotator is not None:
                    scored = self._gene_knowledge_annotator.boost_scores(scored)

                yield from acmg_classifier.classify(scored)
        finally:
            prioritization_engine.close()

    def _build_nmd_lookup(self) -> object | None:
        """Build the NMD-escape transcript lookup from the GTF annotation.

        PVS1 NMD-escape downgrades need only CDS exon coordinates, which the
        GTF supplies without a reference FASTA. Returns None when no GTF is
        configured (API mode, or annotation disabled), in which case PVS1
        keeps its current strength and records transcript_structure missing.
        The index is built once and cached on the instance.
        """
        if self._nmd_lookup_cached is not False:
            return self._nmd_lookup_cached
        lookup: object | None = None
        annotation = self._config.annotation
        if annotation is not None and annotation.gene_annotation_path is not None:
            try:
                from vartriage.annotation.transcript_index import TranscriptCDSIndex

                lookup = TranscriptCDSIndex.build_from_gtf(
                    str(annotation.gene_annotation_path)
                )
            except (OSError, ValueError) as exc:
                logger.warning("NMD lookup unavailable: %s", exc)
                lookup = None
        self._nmd_lookup_cached = lookup
        return lookup

    def _build_acmg_classifier(self) -> ACMGClassifier:
        """Construct the ACMG classifier with the v0.19.0 refinements wired in.

        Supplies the NMD-escape lookup (from the GTF) and the disease-threshold
        flag (from config) so PVS1 NMD downgrades and disease-specific
        PM2/BA1/BS1 thresholds are active in normal pipeline runs, not only when
        a caller constructs the classifier directly.
        """
        return ACMGClassifier(
            nmd_lookup=self._build_nmd_lookup(),  # type: ignore[arg-type]
            use_disease_thresholds=self._config.use_disease_thresholds,
        )

    def _inject_remote_gnomad(self, annotation_engine: AnnotationEngine | None) -> None:
        """Inject remote gnomAD tabix backend if configured and needed."""
        if (
            annotation_engine is not None
            and self._config.remote is not None
            and self._config.remote.is_gnomad_active
            and not annotation_engine.has_frequency_db
        ):
            from vartriage.remote.gnomad import RemoteTabixGnomAD

            annotation_engine.set_frequency_db(RemoteTabixGnomAD(self._config.remote))
            logger.info(
                "Remote gnomAD tabix backend active: %s",
                self._config.remote.gnomad_remote_url,
            )

    def run_with_sv(self) -> None:
        """Execute the main pipeline and SV triage when sv_vcf_path is set.

        Runs the standard SNV pipeline via run(), then if sv_vcf_path is
        configured, also runs the SV triage pipeline and writes a combined
        findings file alongside the main report.
        """
        self.run()

        sv_path = self._config.sv_vcf_path
        if sv_path is None or not sv_path.exists():
            return

        from vartriage.structural.config import SVTriageConfig
        from vartriage.structural.pipeline import SVTriagePipeline

        # Derive SV output path from the main output
        main_output = self._config.output_path
        sv_output = main_output.parent / f"{main_output.stem}_sv{main_output.suffix}"

        sv_config = SVTriageConfig(
            vcf_path=sv_path,
            output_path=sv_output,
            gene_annotation_path=(
                self._config.annotation.gene_annotation_path
                if self._config.annotation is not None
                else None
            ),
            output_format=(
                "json"
                if self._config.report.output_format in ("json", "vcf", "pdf")
                else "csv"
            ),
        )

        sv_pipeline = SVTriagePipeline(sv_config)
        sv_pipeline.run()

        logger.info("SV triage report written to: %s", sv_output)

    def _check_reference_checksums(self) -> None:
        """Log reference file checksums using AuditTrailWriter.

        Delegates SHA-256 computation to
        ``AuditTrailWriter.compute_file_checksum`` so there is a
        single checksum implementation across the codebase.
        """
        from vartriage.reporting.clinical.audit import AuditTrailWriter

        audit_writer = AuditTrailWriter()

        for ref_path in self._collect_reference_paths():
            if not ref_path.exists():
                continue
            try:
                checksum = audit_writer.compute_file_checksum(ref_path)
                logger.debug(
                    "Reference file %s checksum: %s",
                    ref_path,
                    checksum,
                )
            except OSError as exc:
                logger.warning(
                    "Could not compute checksum for %s: %s",
                    ref_path,
                    exc,
                )

    def _collect_reference_paths(self) -> list[Path]:
        """Collect all configured reference file paths.

        Returns
        -------
        list[Path]
            Paths to annotation and prioritization reference files
            that are configured (non-None).
        """
        ref_paths: list[Path] = []
        if self._config.annotation is not None:
            ref_paths.append(self._config.annotation.gene_annotation_path)
            if self._config.annotation.gnomad_path is not None:
                ref_paths.append(self._config.annotation.gnomad_path)
            if self._config.annotation.clinvar_path is not None:
                ref_paths.append(self._config.annotation.clinvar_path)
        pri = self._config.prioritization
        if pri.cadd_scores_path is not None:
            ref_paths.append(pri.cadd_scores_path)
        if pri.revel_scores_path is not None:
            ref_paths.append(pri.revel_scores_path)
        if pri.spliceai_scores_path is not None:
            ref_paths.append(pri.spliceai_scores_path)
        if pri.spliceai_db_path is not None:
            ref_paths.append(pri.spliceai_db_path)
        return ref_paths

    def _compute_reference_checksums(self) -> dict[str, str]:
        """Compute SHA-256 checksums for all reference files.

        Uses ``AuditTrailWriter.compute_file_checksum`` as the
        single checksum implementation. Files that do not exist or
        cannot be read are skipped with a warning.

        Returns
        -------
        dict[str, str]
            Mapping of file path strings to SHA-256 hex digests.
        """
        from vartriage.reporting.clinical.audit import AuditTrailWriter

        audit_writer = AuditTrailWriter()
        checksums: dict[str, str] = {}

        for ref_path in self._collect_reference_paths():
            if not ref_path.exists():
                continue
            try:
                checksums[str(ref_path)] = audit_writer.compute_file_checksum(ref_path)
            except OSError as exc:
                logger.warning(
                    "Could not compute checksum for %s: %s",
                    ref_path,
                    exc,
                )

        return checksums

    def _build_annotated_stream(
        self,
        parser: VCFParser | Iterator[Variant],
        quality_filter: QualityFilter,
        annotation_engine: AnnotationEngine | None,
        api_annotation_engine: object = None,
        sample_names: list[str] | None = None,
    ) -> Iterator[AnnotatedVariant]:
        """Build the filtered and annotated variant stream."""
        stream: Iterator[Variant] = iter(parser)

        # Sample extraction or inheritance (mutually exclusive)
        if self._config.inheritance is not None:
            stream = self._apply_inheritance_filter(stream, parser, sample_names)
        elif self._config.sample is not None:
            stream = self._apply_sample_extraction(stream, parser, sample_names)

        # Region filter (optional, runs before quality filter)
        if self._config.region_filter is not None:
            from vartriage.filter.region_filter import RegionFilter

            region_filter = RegionFilter(self._config.region_filter)
            stream = region_filter.apply(stream)

        filtered = quality_filter.apply(stream)
        if annotation_engine is not None:
            return annotation_engine.annotate(filtered)
        if api_annotation_engine is not None:
            return api_annotation_engine.annotate(filtered)  # type: ignore[attr-defined,no-any-return]
        return self._passthrough_annotation(filtered)

    def _resolve_sample_names(
        self,
        parser: VCFParser | Iterator[Variant],
        sample_names: list[str] | None,
    ) -> list[str]:
        """Resolve sample names from parameter or parser, raising on failure."""
        names = sample_names
        if names is None and hasattr(parser, "sample_names"):
            names = parser.sample_names
        if not names:
            raise ValueError(
                "Sample names required but none available. Ensure the VCF "
                "contains a sample column or pass sample_names explicitly."
            )
        return names

    def _apply_inheritance_filter(
        self,
        stream: Iterator[Variant],
        parser: VCFParser | Iterator[Variant],
        sample_names: list[str] | None,
    ) -> Iterator[Variant]:
        """Apply trio inheritance filtering to the variant stream."""
        from vartriage.filter.inheritance_filter import InheritanceFilter

        names = self._resolve_sample_names(parser, sample_names)
        inheritance_filter = InheritanceFilter(
            self._config.inheritance,  # type: ignore[arg-type]
            names,
        )
        compound_het_active = (
            "compound_het" in self._config.inheritance.patterns  # type: ignore[union-attr]
        )

        if compound_het_active and self._config.annotation is not None:
            return self._apply_compound_het_path(stream, inheritance_filter)

        return inheritance_filter.apply(stream)

    def _apply_compound_het_path(
        self,
        stream: Iterator[Variant],
        inheritance_filter: object,
    ) -> Iterator[Variant]:
        """Handle compound het: annotate first, then filter, then re-attach."""
        from vartriage.annotation.engine import AnnotationEngine

        annotation_engine = AnnotationEngine(self._config.annotation)  # type: ignore[arg-type]
        quality_filter = QualityFilter(self._config.quality_filter)
        filtered = quality_filter.apply(stream)
        annotated_iter = annotation_engine.annotate(filtered)

        annotated_list = list(annotated_iter)
        variants_with_genes = list(self._variants_with_gene_info(iter(annotated_list)))
        inherited_variants = list(
            inheritance_filter.apply(iter(variants_with_genes))  # type: ignore[attr-defined]
        )
        reattached = self._reattach_annotations(inherited_variants, annotated_list)
        return iter(reattached)  # type: ignore[arg-type]

    def _apply_sample_extraction(
        self,
        stream: Iterator[Variant],
        parser: VCFParser | Iterator[Variant],
        sample_names: list[str] | None,
    ) -> Iterator[Variant]:
        """Apply single-sample extraction to the variant stream."""
        from vartriage.filter.sample_extractor import SampleExtractor

        names = self._resolve_sample_names(parser, sample_names)
        sample_extractor = SampleExtractor(
            self._config.sample,  # type: ignore[arg-type]
            names,
        )
        return sample_extractor.apply(stream)

    def _build_api_annotation_engine(self) -> object:
        """Construct the API annotation engine from config.

        Deferred import to avoid pulling httpx into the core module.
        """
        from vartriage.api.annotation_engine import APIAnnotationEngine

        return APIAnnotationEngine(self._config.api)  # type: ignore[arg-type]

    def _validate_config(self, config: PipelineConfig) -> None:
        """Validate all configuration at construction time (fail-fast).

        Checks that reference file paths exist and all sub-config
        parameters are within valid ranges.

        Parameters
        ----------
        config : PipelineConfig
            Configuration to validate.

        Raises
        ------
        ValueError
            If any parameter is out of range.
        FileNotFoundError
            If any required reference file does not exist.
        """
        if config.annotation is not None:
            self._validate_annotation_config(config.annotation)

        self._validate_prioritization_config(config.prioritization)

        if config.gene_filter is not None:
            self._check_path(
                config.gene_filter.gene_list_path,
                "Gene list file",
            )

        if config.region_filter is not None:
            self._check_path(
                config.region_filter.bed_path,
                "BED file",
            )

        if (
            config.inheritance is not None
            and "compound_het" in config.inheritance.patterns
            and config.annotation is None
        ):
            raise ValueError(
                "compound_het pattern requires annotation "
                "configuration (gene annotation reference). "
                "Either provide --gene-annotation and --gnomad, "
                "or remove compound_het from the patterns list."
            )

    def _validate_annotation_config(
        self,
        ann_config: AnnotationConfig,
    ) -> None:
        """Validate annotation reference file paths exist."""
        self._check_path(ann_config.gene_annotation_path, "Gene annotation file")
        if ann_config.gnomad_path is not None:
            self._check_path(ann_config.gnomad_path, "gnomAD reference file")
        if ann_config.clinvar_path is not None:
            self._check_path(ann_config.clinvar_path, "ClinVar reference file")

    def _validate_prioritization_config(
        self,
        pri_config: PrioritizationConfig,
    ) -> None:
        """Validate prioritization score file paths exist."""
        if pri_config.cadd_scores_path is not None:
            self._check_path(pri_config.cadd_scores_path, "CADD scores file")
        if pri_config.revel_scores_path is not None:
            self._check_path(pri_config.revel_scores_path, "REVEL scores file")
        if pri_config.spliceai_scores_path is not None:
            self._check_path(pri_config.spliceai_scores_path, "SpliceAI scores file")
        if pri_config.spliceai_db_path is not None:
            self._check_path(pri_config.spliceai_db_path, "SpliceAI SQLite database")

    def _reattach_annotations(
        self,
        inherited_variants: list[Variant],
        annotated_list: list[AnnotatedVariant],
    ) -> list[AnnotatedVariant]:
        """Re-attach annotation data to variants after inheritance filtering.

        Builds a coordinate lookup from the original annotated variants
        and matches each inherited variant back to its annotation.
        Variants that pass inheritance filtering but have no annotation
        match get INTERGENIC/null as fallback.

        Parameters
        ----------
        inherited_variants : list[Variant]
            Variants that passed inheritance filtering (with
            inheritance_pattern in info).
        annotated_list : list[AnnotatedVariant]
            Original annotated variants before inheritance filtering.

        Returns
        -------
        list[AnnotatedVariant]
            Annotated variants with inheritance metadata preserved
            in the underlying variant's info dict.
        """
        from vartriage.models.variant import AnnotatedVariant, FunctionalConsequence

        # Build lookup by (chrom, pos, ref, alt)
        ann_lookup: dict[tuple[str, int, str, str], AnnotatedVariant] = {}
        for av in annotated_list:
            key = (av.variant.chrom, av.variant.pos, av.variant.ref, av.variant.alt)
            ann_lookup[key] = av

        results: list[AnnotatedVariant] = []
        for v in inherited_variants:
            key = (v.chrom, v.pos, v.ref, v.alt)
            original = ann_lookup.get(key)
            if original is not None:
                # Preserve original annotation, attach inheritance
                # info by replacing the variant's info dict
                results.append(
                    AnnotatedVariant(
                        variant=v,
                        consequence=original.consequence,
                        allele_frequency=original.allele_frequency,
                        clinvar_assertion=original.clinvar_assertion,
                        frequency_unknown=original.frequency_unknown,
                        clinvar_unknown=original.clinvar_unknown,
                        gene_name=original.gene_name,
                    )
                )
            else:
                results.append(
                    AnnotatedVariant(
                        variant=v,
                        consequence=FunctionalConsequence.INTERGENIC,
                        frequency_unknown=True,
                        clinvar_unknown=True,
                    )
                )
        return results

    def _variants_with_gene_info(
        self, annotated: Iterator[AnnotatedVariant]
    ) -> Iterator[Variant]:
        """Extract Variant objects with gene_name copied into info dict.

        Used when compound_het needs gene grouping from annotation
        data. Copies gene_name into info["gene"] so the
        InheritanceFilter can group variants by gene.

        Parameters
        ----------
        annotated : Iterator[AnnotatedVariant]
            Annotated variant stream.

        Yields
        ------
        Variant
            Raw variants with gene info attached.
        """
        for av in annotated:
            v = av.variant
            new_info = dict(v.info)
            if av.gene_name is not None:
                new_info["gene"] = av.gene_name
            yield Variant(
                chrom=v.chrom,
                pos=v.pos,
                id=v.id,
                ref=v.ref,
                alt=v.alt,
                qual=v.qual,
                filter_status=v.filter_status,
                info=new_info,
            )

    def _passthrough_annotation(
        self, variants: Iterator[Variant]
    ) -> Iterator[AnnotatedVariant]:
        """Create AnnotatedVariant wrappers when no annotation config exists.

        Used when the pipeline is run without annotation references. Each
        variant gets an Intergenic consequence, null frequency, and null
        ClinVar assertion, allowing downstream stages to function.

        Parameters
        ----------
        variants : Iterator[Variant]
            Filtered variant stream.

        Yields
        ------
        AnnotatedVariant
            Minimally annotated variant records.
        """
        from vartriage.models.variant import AnnotatedVariant, FunctionalConsequence

        for variant in variants:
            yield AnnotatedVariant(
                variant=variant,
                consequence=FunctionalConsequence.INTERGENIC,
                allele_frequency=None,
                clinvar_assertion=None,
                frequency_unknown=True,
                clinvar_unknown=True,
            )

    @staticmethod
    def _check_path(path: Path, label: str) -> None:
        """Verify a file path exists, raising FileNotFoundError if not.

        Resolves the path to its canonical form before checking, which
        eliminates directory traversal sequences and follows symlinks.

        Parameters
        ----------
        path : Path
            The filesystem path to validate.
        label : str
            Human-readable description used in the error message.

        Raises
        ------
        FileNotFoundError
            If the path does not exist.
        """
        resolved = resolve_path(path)
        if not resolved.exists():
            raise FileNotFoundError(f"{label} not found: {path}")

mito_results property

Access mitochondrial classification results from the last run.

Returns None if mitochondrial analysis was disabled or no chrM variants were present in the input.

qc_report property

Access the QC report from the last run. None if QC was skipped.

warning_accumulator property

Access the warning accumulator tracking MissingDataWarnings.

Returns

WarningAccumulator The shared warning accumulator for the current pipeline run.

run(vcf_path=None, output_path=None)

Execute the full variant prioritization pipeline.

Wires stages sequentially: VCFParser → QualityFilter → AnnotationEngine → PrioritizationEngine → ACMGClassifier → ReportGenerator. Processes data in batches to maintain peak RSS below 2GB for whole-genome scale files (4M+ variants).

Parameters

vcf_path : Path, optional Path to the input VCF file (.vcf or .vcf.gz). If None, uses the path from config. output_path : Path, optional Path for the output report. If None, uses the path from config.

Returns

Path Path to the generated output report file.

Raises

FileNotFoundError If the VCF file or any reference file is missing. ParseError If the VCF file has malformed headers or data lines. IOError If report generation fails due to a write error.

Source code in vartriage/pipeline.py
def run(
    self, vcf_path: Path | None = None, output_path: Path | None = None
) -> Path:
    """Execute the full variant prioritization pipeline.

    Wires stages sequentially: VCFParser → QualityFilter →
    AnnotationEngine → PrioritizationEngine → ACMGClassifier →
    ReportGenerator. Processes data in batches to maintain peak RSS
    below 2GB for whole-genome scale files (4M+ variants).

    Parameters
    ----------
    vcf_path : Path, optional
        Path to the input VCF file (.vcf or .vcf.gz). If None, uses
        the path from config.
    output_path : Path, optional
        Path for the output report. If None, uses the path from config.

    Returns
    -------
    Path
        Path to the generated output report file.

    Raises
    ------
    FileNotFoundError
        If the VCF file or any reference file is missing.
    ParseError
        If the VCF file has malformed headers or data lines.
    IOError
        If report generation fails due to a write error.
    """
    effective_vcf_path = vcf_path or self._config.vcf_path
    effective_output_path = output_path or self._config.output_path

    self._warning_accumulator.reset()

    logger.info(
        "Starting pipeline run: %s → %s", effective_vcf_path, effective_output_path
    )

    # QC pre-flight pass (reads VCF independently before annotation).
    # Result is stored on self._qc_report for report generation and
    # programmatic access via the qc_report property.
    self._run_qc_pass(effective_vcf_path)

    if self._config.clinical_report is not None:
        self._check_reference_checksums()

    (
        quality_filter,
        annotation_engine,
        api_annotation_engine,
        prioritization_engine,
        acmg_classifier,
        report_generator,
    ) = self._build_stages()

    # Per-sample AD/AF are needed for heteroplasmy and maternal
    # inheritance, so extract them whenever mito analysis is active, not
    # only when an inheritance or single-sample selection is configured.
    extract_samples = (
        self._config.inheritance is not None
        or self._config.sample is not None
        or self._mito_enabled
    )

    with VCFParser(
        effective_vcf_path,
        extract_samples=extract_samples,
    ) as parser:
        mito_variants, nuclear_iter = self._split_mito_variants(parser)

        annotated = self._build_annotated_stream(
            nuclear_iter,
            quality_filter,
            annotation_engine,
            api_annotation_engine,
            sample_names=parser.sample_names,
        )

        if self._config.gene_filter is not None:
            from vartriage.filter.gene_filter import GeneFilter

            annotated = GeneFilter(self._config.gene_filter).apply(annotated)

        if self._gene_knowledge_annotator is not None:
            annotated = self._gene_knowledge_annotator.annotate(annotated)

        scored = prioritization_engine.prioritize(annotated)

        if self._gene_knowledge_annotator is not None:
            scored = self._gene_knowledge_annotator.boost_scores(scored)

        classified = acmg_classifier.classify(scored)

        try:
            if self._mito_enabled:
                classified_list = list(classified)
                mito_results = self._run_mito_pipeline(mito_variants)
                result_path = self._generate_report(
                    report_generator,
                    iter(classified_list),
                    effective_output_path,
                    effective_vcf_path,
                    mito_results=mito_results,
                )
            else:
                result_path = self._generate_report(
                    report_generator,
                    classified,
                    effective_output_path,
                    effective_vcf_path,
                    mito_results=None,
                )

            if annotation_engine is not None:
                self._warning_accumulator.add_batch(annotation_engine.warnings)
        finally:
            prioritization_engine.close()

    logger.info(
        "Pipeline completed. Missing data warnings: %d",
        self._warning_accumulator.total_count,
    )
    logger.info("Report written to: %s", result_path)
    return result_path

run_to_classification(vcf_path=None)

Execute the pipeline through classification without report generation.

Runs VCF parsing, quality filtering, annotation, prioritization, and ACMG classification, yielding ClassifiedVariant objects. Skips report writing entirely. Used by CohortPipeline to collect per-sample results without generating throw-away files.

Stage limitations vs full run(): - SampleExtractor is not applied (cohort mode processes each VCF as a single-sample file already). - InheritanceFilter is not applied (trio analysis is a single-proband concern, not meaningful per-sample in a cohort where each VCF represents one individual). - SecondaryFindingsFilter is not applied (cohort mode aggregates by variant, not by clinical reporting context).

These stages are intentionally excluded because cohort analysis operates on pre-separated per-sample VCFs. If your workflow requires sample extraction or inheritance filtering, run the full single-sample pipeline first and feed ClassifiedVariant results to CohortAggregator directly.

Parameters

vcf_path : Path, optional Path to the input VCF file. If None, uses the path from config.

Yields

ClassifiedVariant Classified variants from the pipeline stages.

Source code in vartriage/pipeline.py
def run_to_classification(
    self, vcf_path: Path | None = None
) -> Iterator[ClassifiedVariant]:
    """Execute the pipeline through classification without report generation.

    Runs VCF parsing, quality filtering, annotation, prioritization,
    and ACMG classification, yielding ClassifiedVariant objects. Skips
    report writing entirely. Used by CohortPipeline to collect per-sample
    results without generating throw-away files.

    Stage limitations vs full run():
        - SampleExtractor is not applied (cohort mode processes
          each VCF as a single-sample file already).
        - InheritanceFilter is not applied (trio analysis is a
          single-proband concern, not meaningful per-sample in a
          cohort where each VCF represents one individual).
        - SecondaryFindingsFilter is not applied (cohort mode
          aggregates by variant, not by clinical reporting context).

    These stages are intentionally excluded because cohort analysis
    operates on pre-separated per-sample VCFs. If your workflow
    requires sample extraction or inheritance filtering, run the
    full single-sample pipeline first and feed ClassifiedVariant
    results to CohortAggregator directly.

    Parameters
    ----------
    vcf_path : Path, optional
        Path to the input VCF file. If None, uses the path from config.

    Yields
    ------
    ClassifiedVariant
        Classified variants from the pipeline stages.
    """
    effective_vcf_path = vcf_path or self._config.vcf_path

    quality_filter = QualityFilter(self._config.quality_filter)

    annotation_engine: AnnotationEngine | None = None
    if self._config.annotation is not None:
        annotation_engine = AnnotationEngine(self._config.annotation)

    # Inject remote gnomAD for cohort mode
    self._inject_remote_gnomad(annotation_engine)

    prioritization_engine = PrioritizationEngine(
        self._config.prioritization, remote_config=self._config.remote
    )
    acmg_classifier = self._build_acmg_classifier()

    try:
        with VCFParser(effective_vcf_path) as parser:
            stream: Iterator[Variant] = iter(parser)

            if self._config.region_filter is not None:
                from vartriage.filter.region_filter import RegionFilter

                region_filter = RegionFilter(self._config.region_filter)
                stream = region_filter.apply(stream)

            filtered = quality_filter.apply(stream)

            if annotation_engine is not None:
                annotated = annotation_engine.annotate(filtered)
            else:
                annotated = self._passthrough_annotation(filtered)

            if self._config.gene_filter is not None:
                from vartriage.filter.gene_filter import GeneFilter

                gene_filter = GeneFilter(self._config.gene_filter)
                annotated = gene_filter.apply(annotated)

            # Gene-disease linkage in run_to_classification path
            if self._gene_knowledge_annotator is not None:
                annotated = self._gene_knowledge_annotator.annotate(annotated)

            scored = prioritization_engine.prioritize(annotated)

            # Apply phenotype boost in classification path too
            if self._gene_knowledge_annotator is not None:
                scored = self._gene_knowledge_annotator.boost_scores(scored)

            yield from acmg_classifier.classify(scored)
    finally:
        prioritization_engine.close()

run_with_sv()

Execute the main pipeline and SV triage when sv_vcf_path is set.

Runs the standard SNV pipeline via run(), then if sv_vcf_path is configured, also runs the SV triage pipeline and writes a combined findings file alongside the main report.

Source code in vartriage/pipeline.py
def run_with_sv(self) -> None:
    """Execute the main pipeline and SV triage when sv_vcf_path is set.

    Runs the standard SNV pipeline via run(), then if sv_vcf_path is
    configured, also runs the SV triage pipeline and writes a combined
    findings file alongside the main report.
    """
    self.run()

    sv_path = self._config.sv_vcf_path
    if sv_path is None or not sv_path.exists():
        return

    from vartriage.structural.config import SVTriageConfig
    from vartriage.structural.pipeline import SVTriagePipeline

    # Derive SV output path from the main output
    main_output = self._config.output_path
    sv_output = main_output.parent / f"{main_output.stem}_sv{main_output.suffix}"

    sv_config = SVTriageConfig(
        vcf_path=sv_path,
        output_path=sv_output,
        gene_annotation_path=(
            self._config.annotation.gene_annotation_path
            if self._config.annotation is not None
            else None
        ),
        output_format=(
            "json"
            if self._config.report.output_format in ("json", "vcf", "pdf")
            else "csv"
        ),
    )

    sv_pipeline = SVTriagePipeline(sv_config)
    sv_pipeline.run()

    logger.info("SV triage report written to: %s", sv_output)