Skip to content
Draft
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
46 changes: 38 additions & 8 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,8 @@ on:
branches:
- master
schedule:
- cron: '0 0 1 * *'
- cron: '0 6 * * *'
workflow_dispatch:

jobs:
test:
Expand All @@ -17,25 +18,54 @@ jobs:

steps:
- name: Checkout
uses: actions/checkout@v4
uses: actions/checkout@v7

- name: Setup Micromamba
uses: mamba-org/setup-micromamba@v1
uses: mamba-org/setup-micromamba@v3
with:
micromamba-version: '1.5.10-0' # any version from https://github.com/mamba-org/micromamba-releases
post-cleanup: 'all'

# Pinned, not 'latest': main.nf needs `nextflow.preview.topic`, which 25.x removed once
# topic channels became stable. 24.10.5 is what the cluster runs, so CI matches production.
- name: Setup Nextflow
uses: nf-core/setup-nextflow@v2
- name: Setup Nextflow (minimum supported version)
uses: nf-core/setup-nextflow@v3
with:
version: '25.04.0'

- name: Install nf-test
uses: nf-core/setup-nf-test@v2
with:
version: 0.9.5

- name: Run Tests
run: nf-test test --verbose

test-latest-nextflow:
# Runs on a schedule (and manually), not on every push/PR: this catches a future
# Nextflow release breaking us, without blocking normal PRs on an upstream change
# unrelated to them.
if: github.event_name == 'schedule' || github.event_name == 'workflow_dispatch'
runs-on: ubuntu-latest
timeout-minutes: 60

steps:
- name: Checkout
uses: actions/checkout@v7

- name: Setup Micromamba
uses: mamba-org/setup-micromamba@v3
with:
micromamba-version: '1.5.10-0' # any version from https://github.com/mamba-org/micromamba-releases
post-cleanup: 'all'

- name: Setup Nextflow (latest)
uses: nf-core/setup-nextflow@v3
with:
version: "24.10.5"

# Pinned: nf-test serialises the workflow.trace map in a version-dependent key order, so
# every tests/*.snap must be generated with the same version this installs.
- name: Install nf-test
uses: nf-core/setup-nf-test@v1
uses: nf-core/setup-nf-test@v2
with:
version: 0.9.5

Expand Down
2 changes: 1 addition & 1 deletion .gitignore
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
.nextflow*
work
output
*output/
ngs-agg_revision.yaml
**/.DS_Store
.vscode/
Expand Down
40 changes: 40 additions & 0 deletions BENCHMARK_RESULTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
# Benchmark results: bamadap + fgumi zipper + picodup vs the current pipeline

Performance changes in branch: `bamslice_zip`

**Bottom line**: ~2.4x faster end-to-end on a clean dedicated node (16m/6m45s), with indistinguishable output. Correctness verified in depth (§2);

## 1. What changed, and the measured impact

| Change | Impact |
| --- | --- |
| `fastp`+`bwa`+Picard-dedup pipeline → `bamadap` + `fgumi zipper` + `picodup`, fused into single-pipe `trimAndAlign` + streamed `mergeAndPicodup` | End-to-end: **16m → 6m45s (~2.4x)** |
| Dedup: Picard `MarkDuplicates` → `picodup`, streamed straight off `samtools merge` (no intermediate merged BAM) | **2m02s → 27s (~4.6x)**, peak RSS 8.5GB → <1GB |
| Fixed a real `bwameth.py` bug (naive interleave-detection failed on `bamadap`'s `read/1`,`read/2` mate names, silently mis-converting bisulfite reads); fixed via `bamadap --no-mate-suffix` rather than a downstream patch | Removes a 3x slowdown *and* a silent correctness bug in the align stage |
| `convert_methylkit_to_bed`: single-threaded `gawk` → `mawk` under `parallel --pipepart` | **7m53s → 35s (~14x)** — biggest single bottleneck in either pipeline; benefits `master` too |
| `tasmanian` 1.x (single-threaded, capped at 2M reads, errored on some real reads) → `tasmanian-mismatch` 2.x (parallel, full library) | **3m42s → 17.5s (~13x), on more data** |
| `fastqc` → `falco` | 57.8s → 30.8s (~2x) |
| Fixed `gc_bias` (silently broken: `CollectGcBiasMetrics` requires an `Rscript` on PATH even when no chart is produced; `picard-slim` excludes R) | No-op `Rscript` shim; restores a previously-silent-failing QC step |
| Right-sized cpu reservations for QC steps that don't read `task.cpus` (new `single_threaded_qc` label, fixed at 2 cpus instead of scaling with `--max_cpus`) | Frees queue slots on shared executors; no effect on single-task time |
| Kept `bwameth --read-group` (populated from the uBAM's first `@RG`) instead of dropping it as originally planned; `zipper` still fixes the per-read `RG:Z:` tag | Avoids `picodup` mislabeling every metric row "Unknown Library" |
| Fixed `bamadap`'s EL8 build (glibc-2.34-only symbols, `target-cpu=native` SIGILL'd on other CPUs) | Portability now solved at deploy time, not build time: Capistrano (`capistrano-rust-buildcache`) compiles a `-C target-cpu=native` binary per distinct machine type on first run and caches it, rather than shipping one generic binary |

Chunk-size sweep on the real SGE cluster (18/37/55/90MB): smaller chunks buy wall-clock speed at
the cost of more aggregate CPU-hours (7m14s/7.1 CPU-h at 18MB vs. 10m18s/4.4 CPU-h at 90MB); no
change made to the 37MB default — it's a reasonable middle ground, not a correctness question.

## 2. Correctness (final `.md.bam`, same test uBAM)

| | Baseline | New | Δ |
| --- | --- | --- | --- |
| Total reads | 10,017,473 | 10,017,757 | +0.003% |
| Duplicates | 1,362,823 | 1,362,959 | +0.01% |
| PERCENT_DUPLICATION | 0.136415 | 0.136426 | matches to 4th decimal |
| CpG methylation (Pearson r) | — | **0.999998** | 5 sites (0.0001%) differ by >1pp |

Residual differences trace to expected causes (`bamadap` vs `fastp` trimming, corrected
per-read `RG` moving optical-duplicate grouping, each dedup tool's own library-size extrapolation
formula) — no read loss, no regression. `ngs-aggregate_results` compatibility for
`tasmanian-mismatch` 2.x's new output format was verified against the real deployed parser
(PR #932, merged and deployed 2026-09-07) and an end-to-end real (non-stubbed) aggregation run
succeeded.
9 changes: 3 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -206,9 +206,6 @@ nf-test test --updateSnapshot
```

## Upgrade
Nextflow dropped upstream support for v24 and older in July 2026, but this pipeline still targets
24.10.x -- that is what CI pins and what the reference runs use -- so *main.nf* ships with
`nextflow.preview.topic = true` as its top line.

Nextflow 25.x made topic channels stable and removed that directive. To run on 25.x or newer,
comment out or delete that top line of *main.nf*.
This pipeline requires Nextflow >=25.04, where topic channels (used throughout for
version reporting) are a stable feature rather than a preview one. `main.nf` checks
this at startup and fails immediately with a clear message on older versions.
16 changes: 16 additions & 0 deletions conf/base.config
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,14 @@ process {
cpus = { params.max_cpus ?: 8 }
}

withLabel: dedup {
// Deliberately not capped by params.max_cpus: this label exists so the dedup stage
// can be tuned independently of the trim/align concurrency (see BENCHMARK_PLAN.md phase 6).
// Still capped to whatever the machine actually has, so it doesn't request more cpus
// than any executor can ever grant (e.g. GitHub Actions' 4-core runners).
cpus = { Math.min(16, Runtime.runtime.availableProcessors()) }
}

withLabel: low_cpu {
cpus = 2
memory = { params.max_memory ?: 16.GB }
Expand All @@ -38,6 +46,14 @@ process {
memory = { params.max_memory ?: 8.GB }
}

withLabel: single_threaded_qc {
// For fastqc (falco), gc_bias, insert_size_metrics, picard_metrics, idx_stats:
// none of them read task.cpus, so scaling with params.max_cpus like medium_cpu
// would only reserve idle cpu slots, starving other queued tasks on the local
// executor for no speed benefit.
cpus = 2
}

executor = 'local' // 'sge'. 'slurm', etc...


Expand Down
15 changes: 9 additions & 6 deletions conf/references.config
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,13 @@
// whole-reference curve alone. See "GC bias curves" in the README for the directory's contents.
params {
genomes = [
'T2T_chm13v2.0+bs_controls': [
bwa_index: "/bioinfo/ref/T2T_chm13v2.0+bs_controls/bwameth/T2T_chm13v2.0+bs_controls.fa",
genome_fa: "/bioinfo/ref/T2T_chm13v2.0+bs_controls/bwameth/T2T_chm13v2.0+bs_controls.fa",
genome_fai: "/bioinfo/ref/T2T_chm13v2.0+bs_controls/bwameth/T2T_chm13v2.0+bs_controls.fa.fai",
gc_groups_dir: "/bioinfo/ref/T2T_chm13v2.0+bs_controls/gc_groups",
target_bed: "/bioinfo/ref/t2t_chm13_v2+meth_controls+m13+phix/folded_windows/T2T_chm13v2.0+meth_controls_hairpins.s18w200nogt.bed"
'T2T_chm13v2.0+meth_controls': [
bwa_index: "/bioinfo/ref/t2t_chm13v2+meth_controls/bwameth_index/T2T_chm13v2.0+bs_controls.fa",
genome_fa: "/bioinfo/ref/t2t_chm13v2+meth_controls/bwameth_index/T2T_chm13v2.0+bs_controls.fa",
genome_fai: "/bioinfo/ref/t2t_chm13v2+meth_controls/bwameth_index/T2T_chm13v2.0+bs_controls.fa.fai",
genome_dict:"/bioinfo/ref/t2t_chm13v2+meth_controls/bwameth_index/T2T_chm13v2.0+bs_controls.fa.dict",
gc_groups_dir: "/bioinfo/ref/t2t_chm13v2+meth_controls/gc_groups",
target_bed: "/bioinfo/ref/t2t_chm13v2+meth_controls/folded_windows/T2T_chm13v2.0+meth_controls_hairpins.s18w200nogt.bed",
],
'grcm39+meth_controls': [
bwa_index: "/bioinfo/ref/grcm39+meth_controls/bwameth/grcm39+meth_controls.fa",
Expand Down Expand Up @@ -51,6 +52,7 @@ params {
bwa_index: "${projectDir}/tests/fixtures/reference_files/reference.fa",
genome_fa: "${projectDir}/tests/fixtures/reference_files/reference.fa",
genome_fai: "${projectDir}/tests/fixtures/reference_files/reference.fa.fai",
genome_dict: "${projectDir}/tests/fixtures/reference_files/reference.dict",
gc_groups_dir: "${projectDir}/tests/fixtures/reference_files/gc_groups",
target_bed: "${projectDir}/tests/fixtures/target_bed/emseq_test_regions.bed"
],
Expand All @@ -60,6 +62,7 @@ params {
bwa_index: "${projectDir}/tests/fixtures/reference_files/reference.fa",
genome_fa: "${projectDir}/tests/fixtures/reference_files/reference.fa",
genome_fai: "${projectDir}/tests/fixtures/reference_files/reference.fa.fai",
genome_dict: "${projectDir}/tests/fixtures/reference_files/reference.dict",
target_bed: "${projectDir}/tests/fixtures/target_bed/emseq_test_regions.bed"
]
]
Expand Down
31 changes: 19 additions & 12 deletions conf/test.config
Original file line number Diff line number Diff line change
Expand Up @@ -11,20 +11,27 @@
*/

process{
executor = 'local'
executor = 'local'
}

params {
config_profile_name = 'Test profile'
config_profile_description = 'Minimal test dataset to check pipeline function'

// Limit resources so that this can run on GitHub Actions
max_cpus = 2
max_memory = '6.GB'
max_time = '6.h'

outputDir = "test_output"
}
config_profile_name = 'Test profile'
config_profile_description = 'Minimal test dataset to check pipeline function'

// Limit resources so that this can run on GitHub Actions
max_cpus = 2
max_memory = '6.GB'
max_time = '6.h'

outputDir = "test_output"

// Not yet on bioconda. nf-test's pipeline tests run trimAndAlign/mergeAndPicodup under
// -stub-run (see their stub: blocks), so CI never actually invokes these paths -- they
// only matter for a manual, non-stub `nextflow run main.nf -profile test`.
bamadap_bin = '/appdev/langhorst/src/bamadap/target-el8/release/bamadap'
picodup_bin = '/appdev/langhorst/src/picodup/target-el8/release/picodup'
bamslice_bin = '/appdev/langhorst/src/bamslice/target-el8/release/bamslice'
}

trace.overwrite = true
trace.overwrite = true

57 changes: 31 additions & 26 deletions lib/notifications.nf
Original file line number Diff line number Diff line change
@@ -1,40 +1,41 @@
def notificationHtml(String status) {
// `wf` is the live WorkflowMetadata object captured by registerEmailNotifications
def notificationHtml(String status, String pipelineName, wf) {
def isSuccess = status == 'SUCCESS'
def bannerBg = isSuccess ? '#dff0d8' : '#f2dede'
def bannerBorder = isSuccess ? '#d6e9c6' : '#ebccd1'
def bannerColor = isSuccess ? '#3c763d' : '#a94442'
def bannerMsg = isSuccess ? 'Execution completed successfully!' : 'Execution failed!'

def rows = []
rows << ['Pipeline', params.workflow ?: '-']
rows << ['Run name', workflow.runName]
rows << ['Launch time', workflow.start]
if (workflow.complete) {
rows << ['Ending time', "${workflow.complete} (duration: ${workflow.duration})"]
rows << ['Pipeline', pipelineName ?: '-']
rows << ['Run name', wf.runName]
rows << ['Launch time', wf.start]
if (wf.complete) {
rows << ['Ending time', "${wf.complete} (duration: ${wf.duration})"]
}
def stats = null
try { stats = workflow.stats } catch (Exception e) { }
try { stats = wf.stats } catch (Exception e) { }
if (stats) {
rows << ['Tasks stats', "Succeeded: ${stats.succeededCount} &nbsp; Cached: ${stats.cachedCount} &nbsp; Ignored: ${stats.ignoredCount} &nbsp; Failed: ${stats.failedCount}"]
}
rows << ['Launch directory', workflow.launchDir]
rows << ['Work directory', workflow.workDir]
rows << ['Project directory', workflow.projectDir]
rows << ['Script name', workflow.scriptName]
rows << ['Script ID', workflow.scriptId]
rows << ['Workflow session', workflow.sessionId]
if (workflow.profile) rows << ['Workflow profile', workflow.profile]
rows << ['Nextflow version', "${workflow.nextflow.version}, build ${workflow.nextflow.build}"]
rows << ['Launch directory', wf.launchDir]
rows << ['Work directory', wf.workDir]
rows << ['Project directory', wf.projectDir]
rows << ['Script name', wf.scriptName]
rows << ['Script ID', wf.scriptId]
rows << ['Workflow session', wf.sessionId]
if (wf.profile) rows << ['Workflow profile', wf.profile]
rows << ['Nextflow version', "${wf.nextflow.version}, build ${wf.nextflow.build}"]

def rowsHtml = rows.collect { pair ->
" <tr><td style=\"padding:4px 10px;color:#666;width:180px;vertical-align:top;\">${pair[0]}</td><td style=\"padding:4px 10px;word-break:break-all;\">${pair[1] ?: '-'}</td></tr>"
}.join('\n')

def errorBlock = ''
if (!isSuccess && workflow.errorMessage) {
if (!isSuccess && wf.errorMessage) {
errorBlock = """
<p><b>Error:</b></p>
<pre style="background:#f9f2f2;padding:10px;border-radius:4px;color:#a94442;white-space:pre-wrap;">${workflow.errorMessage}</pre>
<pre style="background:#f9f2f2;padding:10px;border-radius:4px;color:#a94442;white-space:pre-wrap;">${wf.errorMessage}</pre>
"""
}

Expand All @@ -44,14 +45,14 @@ def notificationHtml(String status) {
<head><meta charset="utf-8"></head>
<body style="font-family:Helvetica,Arial,sans-serif;max-width:800px;color:#333;padding:20px;">
<h1 style="border-bottom:1px solid #ddd;padding-bottom:10px;">Workflow ${isSuccess ? 'completion' : 'failure'} notification</h1>
<h2 style="margin-top:0;">Run Name: ${workflow.runName}</h2>
<h2 style="margin-top:0;">Run Name: ${wf.runName}</h2>

<div style="background:${bannerBg};border:1px solid ${bannerBorder};color:${bannerColor};padding:10px 15px;border-radius:4px;margin:15px 0;">
${bannerMsg}
</div>
${errorBlock}
<p>The command used to launch the workflow was as follows:</p>
<pre style="background:#f4f4f4;padding:10px;border-radius:4px;font-size:13px;white-space:pre-wrap;word-break:break-all;">${workflow.commandLine}</pre>
<pre style="background:#f4f4f4;padding:10px;border-radius:4px;font-size:13px;white-space:pre-wrap;word-break:break-all;">${wf.commandLine}</pre>

<h2 style="border-bottom:1px solid #ddd;padding-bottom:10px;margin-top:30px;">Execution summary</h2>
<table style="border-collapse:collapse;font-size:14px;">
Expand All @@ -65,29 +66,33 @@ ${rowsHtml}
def registerEmailNotifications() {
if (params.dry_run || workflow.stubRun) return

// workflow.onError/onComplete run detached from this binding, so bare `params.*`/`workflow.*`
// throw NPEs there; capture into locals below so closure capture picks them up instead.
def wf = workflow
def recipients = [params.email, params.admin_email].findAll { it }.join(',')
def pipelineName = params.workflow ?: 'unknown'

def notified = false

workflow.onError {
def recipients = [params.email, params.admin_email].findAll { it }.join(',')
if (!recipients) return
sendMail(
to: recipients,
subject: "[Pipeline FAILED] ${params.workflow ?: 'unknown'} - ${workflow.runName}",
body: notificationHtml('FAILED'),
subject: "[Pipeline FAILED] ${pipelineName} - ${wf.runName}",
body: notificationHtml('FAILED', pipelineName, wf),
type: 'text/html'
)
notified = true
}

workflow.onComplete {
if (notified) return
def recipients = [params.email, params.admin_email].findAll { it }.join(',')
if (!recipients) return
def status = workflow.success ? 'SUCCESS' : 'FAILED'
def status = wf.success ? 'SUCCESS' : 'FAILED'
sendMail(
to: recipients,
subject: "[Pipeline ${status}] ${params.workflow} - ${workflow.runName}",
body: notificationHtml(status),
subject: "[Pipeline ${status}] ${pipelineName} - ${wf.runName}",
body: notificationHtml(status, pipelineName, wf),
type: 'text/html'
)
}
Expand Down
Loading