use anyhow::{Context, Result};
use serde_json::{Map, Value};
fn required_i64(obj: &Value, field: &str) -> Result<i64> {
let value = obj
.get(field)
.with_context(|| format!("missing required field '{field}'"))?;
value
.as_i64()
.with_context(|| format!("field '{field}' is not an integer (found {value})"))
}
fn required_f64(obj: &Value, field: &str) -> Result<f64> {
let value = obj
.get(field)
.with_context(|| format!("missing required field '{field}'"))?;
value
.as_f64()
.with_context(|| format!("field '{field}' is not a number (found {value})"))
}
fn element_as_f64(value: &Value) -> Result<f64> {
value
.as_f64()
.with_context(|| format!("array contains a non-numeric value ({value})"))
}
pub fn sum_arrays(arrays: &[&[Value]]) -> Result<Vec<Value>> {
if arrays.is_empty() {
return Ok(Vec::new());
}
let max_len = arrays.iter().map(|a| a.len()).max().unwrap_or(0);
let all_integers = arrays
.iter()
.all(|arr| arr.iter().all(|v| v.is_i64() || v.is_u64()));
if all_integers {
let mut result = vec![0_i64; max_len];
for arr in arrays {
for (i, val) in arr.iter().enumerate() {
result[i] = result[i].saturating_add(val.as_i64().unwrap_or(0));
}
}
Ok(result.into_iter().map(Value::from).collect())
} else {
let mut result = vec![0_f64; max_len];
for arr in arrays {
for (i, val) in arr.iter().enumerate() {
result[i] += element_as_f64(val)?;
}
}
Ok(result.into_iter().map(Value::from).collect())
}
}
pub fn weighted_average_arrays(arrays: &[&[Value]], weights: &[f64]) -> Result<Vec<Value>> {
if arrays.is_empty() {
return Ok(Vec::new());
}
let max_len = arrays.iter().map(|a| a.len()).max().unwrap_or(0);
let mut sums = vec![0_f64; max_len];
let mut weight_totals = vec![0_f64; max_len];
for (arr, &weight) in arrays.iter().zip(weights.iter()) {
for (i, val) in arr.iter().enumerate() {
sums[i] = element_as_f64(val)?.mul_add(weight, sums[i]);
weight_totals[i] += weight;
}
}
Ok(sums
.iter()
.zip(weight_totals.iter())
.map(|(&sum, &total)| {
if total > 0.0 {
Value::from(sum / total)
} else {
Value::from(0.0)
}
})
.collect())
}
fn merge_int_counts(dicts: &[&Map<String, Value>], item_label: &str) -> Result<Map<String, Value>> {
let mut merged: Map<String, Value> = Map::new();
for dict in dicts {
for (key, count) in *dict {
let count = count.as_i64().with_context(|| {
format!("{item_label} '{key}' count is not an integer (found {count})")
})?;
let existing = merged.get(key).and_then(Value::as_i64).unwrap_or(0);
merged.insert(key.clone(), Value::from(existing.saturating_add(count)));
}
}
Ok(merged)
}
fn merge_kmer_counts(dicts: &[&Map<String, Value>]) -> Result<Map<String, Value>> {
merge_int_counts(dicts, "kmer")
}
fn merge_curve_maps(curves: &[&Map<String, Value>], weights: &[f64]) -> Result<Map<String, Value>> {
if curves.is_empty() {
return Ok(Map::new());
}
let mut seen = std::collections::HashSet::new();
let all_keys: Vec<&String> = curves
.iter()
.flat_map(|c| c.keys())
.filter(|k| seen.insert(k.as_str()))
.collect();
let mut merged = Map::new();
for key in all_keys {
let mut arrays: Vec<&[Value]> = Vec::new();
let mut arr_weights: Vec<f64> = Vec::new();
for (curve, &weight) in curves.iter().zip(weights.iter()) {
if let Some(value) = curve.get(key) {
let array = value
.as_array()
.with_context(|| format!("curve '{key}' is not an array"))?;
arrays.push(array.as_slice());
arr_weights.push(weight);
}
}
let averaged = weighted_average_arrays(&arrays, &arr_weights)
.with_context(|| format!("averaging curve '{key}'"))?;
merged.insert(key.clone(), Value::Array(averaged));
}
Ok(merged)
}
fn merge_adapter_counts(dicts: &[&Map<String, Value>]) -> Result<Map<String, Value>> {
merge_int_counts(dicts, "adapter")
}
fn merge_filtering_result(results: &[Value]) -> Result<Value> {
const FIELDS: &[&str] = &[
"passed_filter_reads",
"low_quality_reads",
"too_many_N_reads",
"adapter_dimer_reads",
"too_short_reads",
"too_long_reads",
];
let has_dimer = results
.iter()
.any(|r| r.get("adapter_dimer_reads").is_some());
let mut merged = Map::new();
for field in FIELDS {
if *field == "adapter_dimer_reads" && !has_dimer {
continue;
}
let mut sum = 0_i64;
for (i, result) in results.iter().enumerate() {
sum += required_i64(result, field).with_context(|| format!("filtering_result[{i}]"))?;
}
merged.insert((*field).to_string(), Value::from(sum));
}
Ok(Value::Object(merged))
}
fn merge_insert_size(sizes: &[Value]) -> Result<Value> {
let mut histograms: Vec<&[Value]> = Vec::new();
for (i, size) in sizes.iter().enumerate() {
let histogram = size
.get("histogram")
.with_context(|| format!("insert_size[{i}]: missing 'histogram'"))?
.as_array()
.with_context(|| format!("insert_size[{i}]: 'histogram' is not an array"))?;
histograms.push(histogram.as_slice());
}
let merged_histogram = sum_arrays(&histograms).context("insert_size histogram")?;
let mut unknown = 0_i64;
for (i, size) in sizes.iter().enumerate() {
unknown += required_i64(size, "unknown").with_context(|| format!("insert_size[{i}]"))?;
}
let peak = merged_histogram
.iter()
.enumerate()
.max_by(|(_, a), (_, b)| {
a.as_f64()
.partial_cmp(&b.as_f64())
.unwrap_or(std::cmp::Ordering::Equal)
})
.map_or(0, |(i, _)| i);
let mut merged = Map::new();
merged.insert("peak".to_string(), Value::from(peak));
merged.insert("unknown".to_string(), Value::from(unknown));
merged.insert("histogram".to_string(), Value::Array(merged_histogram));
Ok(Value::Object(merged))
}
fn merge_duplication(dups: &[Value], summaries: &[Value]) -> Result<Value> {
let mut weighted_sum = 0.0_f64;
let mut total_reads = 0.0_f64;
for (i, (dup, summary)) in dups.iter().zip(summaries.iter()).enumerate() {
let reads = summary
.get("before_filtering")
.and_then(|b| b.get("total_reads"))
.with_context(|| {
format!("duplication[{i}]: missing summary.before_filtering.total_reads")
})?
.as_f64()
.with_context(|| format!("duplication[{i}]: total_reads is not a number"))?;
let rate = required_f64(dup, "rate").with_context(|| format!("duplication[{i}]"))?;
weighted_sum = rate.mul_add(reads, weighted_sum);
total_reads += reads;
}
let rate = if total_reads > 0.0 {
weighted_sum / total_reads
} else {
0.0
};
let mut merged = Map::new();
merged.insert("rate".to_string(), Value::from(rate));
Ok(Value::Object(merged))
}
#[allow(clippy::cast_possible_truncation)]
fn apply_summary_rates(merged: &mut Map<String, Value>, stats: &[Value]) -> Result<()> {
let total_bases = merged
.get("total_bases")
.and_then(Value::as_f64)
.unwrap_or(0.0);
let total_reads = merged
.get("total_reads")
.and_then(Value::as_f64)
.unwrap_or(0.0);
if total_bases > 0.0 {
let q20 = merged
.get("q20_bases")
.and_then(Value::as_f64)
.unwrap_or(0.0);
let q30 = merged
.get("q30_bases")
.and_then(Value::as_f64)
.unwrap_or(0.0);
merged.insert("q20_rate".to_string(), Value::from(q20 / total_bases));
merged.insert("q30_rate".to_string(), Value::from(q30 / total_bases));
} else {
merged.insert("q20_rate".to_string(), Value::from(0.0));
merged.insert("q30_rate".to_string(), Value::from(0.0));
}
if total_reads > 0.0 {
for field in &["read1_mean_length", "read2_mean_length"] {
if stats.iter().all(|s| s.get(*field).is_none()) {
continue;
}
let mut weighted_sum = 0.0_f64;
for (i, stat) in stats.iter().enumerate() {
let mean = required_f64(stat, field).with_context(|| format!("summary[{i}]"))?;
let reads =
required_f64(stat, "total_reads").with_context(|| format!("summary[{i}]"))?;
weighted_sum = mean.mul_add(reads, weighted_sum);
}
merged.insert(
(*field).to_string(),
Value::from((weighted_sum / total_reads) as i64),
);
}
}
if total_bases > 0.0 && stats.iter().any(|s| s.get("gc_content").is_some()) {
let mut total_gc = 0.0_f64;
for (i, stat) in stats.iter().enumerate() {
let gc = required_f64(stat, "gc_content").with_context(|| format!("summary[{i}]"))?;
let bases =
required_f64(stat, "total_bases").with_context(|| format!("summary[{i}]"))?;
total_gc = gc.mul_add(bases, total_gc);
}
merged.insert(
"gc_content".to_string(),
Value::from(total_gc / total_bases),
);
}
Ok(())
}
#[derive(Clone, Copy)]
pub enum ReadStatsKind {
Summary,
PerRead,
}
pub fn merge_read_stats(stats: &[Value], stage_name: &str, kind: ReadStatsKind) -> Result<Value> {
let mut merged = Map::new();
for field in &["total_reads", "total_bases", "q20_bases", "q30_bases"] {
let mut sum = 0_i64;
for (i, stat) in stats.iter().enumerate() {
sum += required_i64(stat, field).with_context(|| format!("{stage_name}[{i}]"))?;
}
merged.insert((*field).to_string(), Value::from(sum));
}
match kind {
ReadStatsKind::PerRead => {
if let Some(cycles) = stats.first().and_then(|s| s.get("total_cycles")) {
merged.insert("total_cycles".to_string(), cycles.clone());
}
}
ReadStatsKind::Summary => {
apply_summary_rates(&mut merged, stats)?;
}
}
for section in &["quality_curves", "content_curves"] {
let mut maps: Vec<&Map<String, Value>> = Vec::new();
let mut weights: Vec<f64> = Vec::new();
for stat in stats {
if let Some(obj) = stat.get(*section).and_then(Value::as_object) {
maps.push(obj);
weights.push(
stat.get("total_reads")
.and_then(Value::as_f64)
.unwrap_or(0.0),
);
}
}
if !maps.is_empty() {
let curves = merge_curve_maps(&maps, &weights)
.with_context(|| format!("{stage_name} {section}"))?;
merged.insert((*section).to_string(), Value::Object(curves));
}
}
let kmer_counts: Vec<&Map<String, Value>> = stats
.iter()
.filter_map(|s| s.get("kmer_count")?.as_object())
.collect();
if !kmer_counts.is_empty() {
let merged_kmers =
merge_kmer_counts(&kmer_counts).with_context(|| format!("{stage_name} kmer_count"))?;
merged.insert("kmer_count".to_string(), Value::Object(merged_kmers));
}
if stats
.iter()
.any(|s| s.get("overrepresented_sequences").is_some())
{
merged.insert(
"overrepresented_sequences".to_string(),
Value::Object(Map::new()),
);
}
Ok(Value::Object(merged))
}
fn merge_adapter_cutting(cuttings: &[Value]) -> Result<Value> {
let mut merged = Map::new();
let mut trimmed_reads = 0_i64;
let mut trimmed_bases = 0_i64;
for (i, cutting) in cuttings.iter().enumerate() {
trimmed_reads += required_i64(cutting, "adapter_trimmed_reads")
.with_context(|| format!("adapter_cutting[{i}]"))?;
trimmed_bases += required_i64(cutting, "adapter_trimmed_bases")
.with_context(|| format!("adapter_cutting[{i}]"))?;
}
merged.insert(
"adapter_trimmed_reads".to_string(),
Value::from(trimmed_reads),
);
merged.insert(
"adapter_trimmed_bases".to_string(),
Value::from(trimmed_bases),
);
for field in &["read1_adapter_sequence", "read2_adapter_sequence"] {
if let Some(seq) = cuttings.first().and_then(|ac| ac.get(*field)) {
merged.insert((*field).to_string(), seq.clone());
}
}
for field in &["read1_adapter_counts", "read2_adapter_counts"] {
if cuttings.iter().any(|ac| ac.get(*field).is_some()) {
let mut counts: Vec<&Map<String, Value>> = Vec::new();
for (i, cutting) in cuttings.iter().enumerate() {
if let Some(value) = cutting.get(*field) {
let object = value.as_object().with_context(|| {
format!("adapter_cutting[{i}]: '{field}' is not an object")
})?;
counts.push(object);
}
}
let merged_counts = merge_adapter_counts(&counts)
.with_context(|| format!("adapter_cutting {field}"))?;
merged.insert((*field).to_string(), Value::Object(merged_counts));
}
}
Ok(Value::Object(merged))
}
fn require_section(data_list: &[Value], keys: &[&str]) -> Result<Vec<Value>> {
let mut out = Vec::with_capacity(data_list.len());
for (i, data) in data_list.iter().enumerate() {
let mut current = data;
for key in keys {
current = current.get(*key).with_context(|| {
format!("file[{i}]: missing required section '{}'", keys.join("."))
})?;
}
out.push(current.clone());
}
Ok(out)
}
pub fn merge_jsons(data_list: &[Value]) -> Result<Value> {
if data_list.is_empty() {
return Ok(Value::Object(Map::new()));
}
let mut merged = Map::new();
let mut summary = Map::new();
if let Some(first_summary) = data_list[0].get("summary") {
for field in &["fastp_version", "sequencing"] {
if let Some(val) = first_summary.get(*field) {
summary.insert((*field).to_string(), val.clone());
}
}
}
let before = require_section(data_list, &["summary", "before_filtering"])?;
let after = require_section(data_list, &["summary", "after_filtering"])?;
summary.insert(
"before_filtering".to_string(),
merge_read_stats(&before, "before_filtering", ReadStatsKind::Summary)?,
);
summary.insert(
"after_filtering".to_string(),
merge_read_stats(&after, "after_filtering", ReadStatsKind::Summary)?,
);
merged.insert("summary".to_string(), Value::Object(summary));
let filtering_results = require_section(data_list, &["filtering_result"])?;
merged.insert(
"filtering_result".to_string(),
merge_filtering_result(&filtering_results)?,
);
let dups = require_section(data_list, &["duplication"])?;
let summaries = require_section(data_list, &["summary"])?;
merged.insert(
"duplication".to_string(),
merge_duplication(&dups, &summaries)?,
);
if data_list.iter().any(|d| d.get("insert_size").is_some()) {
let insert_sizes = require_section(data_list, &["insert_size"])?;
merged.insert("insert_size".to_string(), merge_insert_size(&insert_sizes)?);
}
let adapter_cuttings = require_section(data_list, &["adapter_cutting"])?;
merged.insert(
"adapter_cutting".to_string(),
merge_adapter_cutting(&adapter_cuttings)?,
);
for section in &[
"read1_before_filtering",
"read2_before_filtering",
"read1_after_filtering",
"read2_after_filtering",
] {
if data_list.iter().any(|d| d.get(*section).is_some()) {
let values = require_section(data_list, &[section])?;
merged.insert(
(*section).to_string(),
merge_read_stats(&values, section, ReadStatsKind::PerRead)?,
);
}
}
Ok(Value::Object(merged))
}
fn strip_trailing_zeros(s: &str) -> &str {
if s.contains('.') {
s.trim_end_matches('0').trim_end_matches('.')
} else {
s
}
}
fn fastp_float_repr(f: f64) -> String {
if f == 0.0 {
return "0".to_string();
}
if f.is_nan() {
return "nan".to_string();
}
if f.is_infinite() {
return if f < 0.0 { "-inf" } else { "inf" }.to_string();
}
let sci = format!("{f:.5e}");
let (mantissa, exp_str) = sci
.split_once('e')
.expect("scientific format always contains 'e'");
let exp: i32 = exp_str.parse().expect("scientific exponent is an integer");
let negative = mantissa.starts_with('-');
let mant = mantissa.trim_start_matches('-');
let digits: String = mant.chars().filter(|&c| c != '.').collect();
let mut out = String::new();
if negative {
out.push('-');
}
if (-4..6).contains(&exp) {
let decpt = exp + 1;
let body = if decpt <= 0 {
let mut s = String::from("0.");
for _ in 0..-decpt {
s.push('0');
}
s.push_str(&digits);
s
} else {
let split = usize::try_from(decpt).unwrap_or(0);
if split >= digits.len() {
let mut s = digits.clone();
for _ in 0..(split - digits.len()) {
s.push('0');
}
s
} else {
format!("{}.{}", &digits[..split], &digits[split..])
}
};
out.push_str(strip_trailing_zeros(&body));
} else {
let rest = digits[1..].trim_end_matches('0');
out.push_str(&digits[..1]);
if !rest.is_empty() {
out.push('.');
out.push_str(rest);
}
out.push('e');
out.push(if exp < 0 { '-' } else { '+' });
let magnitude = exp.abs();
if magnitude < 10 {
out.push('0');
}
out.push_str(&magnitude.to_string());
}
out
}
fn format_number(n: &serde_json::Number) -> String {
if n.is_i64() || n.is_u64() {
n.to_string()
} else if let Some(f) = n.as_f64() {
fastp_float_repr(f)
} else {
n.to_string()
}
}
fn push_indent(out: &mut String, depth: usize) {
for _ in 0..depth {
out.push('\t');
}
}
fn write_array(arr: &[Value], out: &mut String) -> Result<()> {
out.push('[');
for (i, v) in arr.iter().enumerate() {
if i > 0 {
out.push(',');
}
match v {
Value::Number(n) => out.push_str(&format_number(n)),
other => write_value(other, 0, out)?,
}
}
out.push(']');
Ok(())
}
fn write_kmer_count(map: &Map<String, Value>, depth: usize, out: &mut String) -> Result<()> {
if map.is_empty() {
out.push_str("{}");
return Ok(());
}
let entries: Vec<String> = map
.iter()
.map(|(k, v)| {
Ok(format!(
"{}: {}",
serde_json::to_string(k)?,
format_number_value(v)
))
})
.collect::<Result<_>>()?;
out.push_str("{\n");
let mut i = 0;
while i < entries.len() {
let end = (i + 16).min(entries.len());
push_indent(out, depth + 1);
out.push_str(&entries[i..end].join(",\t\t\t"));
if end < entries.len() {
out.push(',');
}
out.push('\n');
i += 16;
}
push_indent(out, depth);
out.push('}');
Ok(())
}
fn format_number_value(v: &Value) -> String {
match v {
Value::Number(n) => format_number(n),
other => other.to_string(),
}
}
fn write_object(map: &Map<String, Value>, depth: usize, out: &mut String) -> Result<()> {
if map.is_empty() {
out.push_str("{}");
return Ok(());
}
out.push_str("{\n");
let n = map.len();
for (i, (k, v)) in map.iter().enumerate() {
push_indent(out, depth + 1);
out.push_str(&serde_json::to_string(k)?);
match v {
Value::Array(arr) => {
out.push(':');
write_array(arr, out)?;
}
Value::Object(obj) if k == "kmer_count" => {
out.push_str(": ");
write_kmer_count(obj, depth + 1, out)?;
}
_ => {
out.push_str(": ");
write_value(v, depth + 1, out)?;
}
}
if i + 1 < n {
out.push(',');
}
out.push('\n');
}
push_indent(out, depth);
out.push('}');
Ok(())
}
fn write_value(value: &Value, depth: usize, out: &mut String) -> Result<()> {
match value {
Value::Object(map) => write_object(map, depth, out),
Value::Array(arr) => write_array(arr, out),
Value::String(_) => {
out.push_str(&serde_json::to_string(value)?);
Ok(())
}
Value::Number(n) => {
out.push_str(&format_number(n));
Ok(())
}
Value::Bool(b) => {
out.push_str(if *b { "true" } else { "false" });
Ok(())
}
Value::Null => {
out.push_str("null");
Ok(())
}
}
}
pub fn format_json(value: &Value) -> Result<String> {
let mut out = String::new();
write_value(value, 0, &mut out)?;
Ok(out)
}