← Files NGS Analysis WorkbenchARCHIVED FILE
mcp/ngs_workbench_execution_monitoring/nextflow_trace.py
7.55 KB · Sep 30, 2026 · 23:20 UTC
"""Nextflow ``trace.txt`` execution evidence adapter."""
from __future__ import annotations
import csv
from collections.abc import Iterable
from pathlib import Path
from typing import Any
from .models import (
ExecutionEvidence,
ExecutionObservation,
ExecutionProgress,
TaskAttempt,
TaskState,
count_attempts,
discover_processes,
empty_observation,
limit_attempts,
)
_TERMINAL_STATES = {"completed", "failed", "cached", "aborted"}
class NextflowTraceObserver:
"""Normalize task attempts that Nextflow has flushed to ``trace.txt``."""
binding = "nextflow"
engine = "nextflow"
evidence_kind = "nextflow_trace"
file_patterns = ("logs/nextflow.trace.txt", "workflow/trace.txt")
def observe(self, run_dir: Path, _workflow_status: str) -> ExecutionObservation:
trace_path = run_dir / self.file_patterns[0]
if not trace_path.is_file():
trace_path = run_dir / self.file_patterns[1]
if not trace_path.is_file():
return empty_observation(
engine=self.engine,
evidence_kind=self.evidence_kind,
evidence_path=str(trace_path),
reason="trace.txt has not been produced yet",
)
try:
with trace_path.open(encoding="utf-8-sig", newline="") as lines:
return self.observe_lines(
lines,
evidence_path=str(trace_path),
workflow_status=_workflow_status,
)
except (OSError, UnicodeError, csv.Error) as exc:
return empty_observation(
engine=self.engine,
evidence_kind=self.evidence_kind,
evidence_path=str(trace_path),
reason=f"trace.txt could not be read: {exc}",
)
def observe_lines(
self,
lines: Iterable[str],
*,
evidence_path: str,
workflow_status: str,
) -> ExecutionObservation:
"""Normalize one Nextflow trace supplied by a local or remote transport."""
del workflow_status
try:
attempts, ignored_records = self._read_attempts(lines)
except csv.Error as exc:
return empty_observation(
engine=self.engine,
evidence_kind=self.evidence_kind,
evidence_path=evidence_path,
reason=f"trace.txt could not be parsed: {exc}",
)
counts = count_attempts(attempts)
finished_attempts = sum(getattr(counts, state) for state in _TERMINAL_STATES)
reason = "the trace contains observed attempts, not a planned task denominator"
if not attempts:
reason = "trace.txt exists but does not contain complete task records yet"
returned_attempts, attempts_truncated = limit_attempts(attempts)
return ExecutionObservation(
engine=self.engine,
evidence=ExecutionEvidence(
kind=self.evidence_kind,
path=evidence_path,
available=True,
append_only=True,
coverage="flushed task attempts",
reason=reason if not attempts else None,
record_count=len(attempts),
ignored_record_count=ignored_records,
),
counts=counts,
progress=ExecutionProgress(
completed=counts.completed + counts.cached,
finished_attempts=finished_attempts,
total=None,
determinate=False,
reason=reason,
),
structure=discover_processes(
attempts,
ordering="first_terminal_observation",
),
attempts_truncated=attempts_truncated,
attempts=returned_attempts,
)
def _read_attempts(self, lines: Iterable[str]) -> tuple[list[TaskAttempt], int]:
attempts: list[TaskAttempt] = []
ignored_records = 0
seen_attempt_ids: set[str] = set()
reader = csv.DictReader(lines, delimiter="\t")
if not reader.fieldnames or not {"name", "status"}.issubset(reader.fieldnames):
return attempts, 1
for sequence, row in enumerate(reader, start=1):
if (
not self._has_record(row)
or not _optional_text(row.get("name"))
or not _optional_text(row.get("status"))
):
ignored_records += 1
continue
attempt = self._parse_attempt(row, sequence)
if attempt.attempt_id in seen_attempt_ids:
attempt.attempt_id = f"{attempt.attempt_id}:{sequence}"
seen_attempt_ids.add(attempt.attempt_id)
attempts.append(attempt)
return attempts, ignored_records
@staticmethod
def _has_record(row: dict[str | None, Any]) -> bool:
return any(value not in {None, ""} for key, value in row.items() if key is not None)
@staticmethod
def _parse_attempt(row: dict[str | None, Any], sequence: int) -> TaskAttempt:
raw_name = _optional_text(row.get("name")) or "unknown"
process_name, tag = _split_task_name(raw_name)
attempt_number = _optional_int(row.get("attempt"))
task_id = _optional_text(row.get("task_id"))
task_hash = _optional_text(row.get("hash"))
base_attempt_id = task_id or task_hash or str(sequence)
attempt_id = (
f"{base_attempt_id}:{attempt_number}"
if attempt_number is not None and attempt_number > 1
else base_attempt_id
)
return TaskAttempt(
attempt_id=attempt_id,
attempt_number=attempt_number,
task_id=task_id,
task_hash=task_hash,
native_id=_optional_text(row.get("native_id")),
process=process_name,
process_label=process_name.rsplit(":", 1)[-1],
sample_or_shard=tag or _optional_text(row.get("tag")),
state=_normalize_state(row.get("status")),
exit_code=_optional_int(row.get("exit")),
submitted_at=_optional_text(row.get("submit")),
started_at=_optional_text(row.get("start")),
completed_at=_optional_text(row.get("complete")),
duration=_optional_text(row.get("duration")),
realtime=_optional_text(row.get("realtime")),
cpu=_optional_text(row.get("%cpu")),
peak_rss=_optional_text(row.get("peak_rss")),
peak_vmem=_optional_text(row.get("peak_vmem")),
workdir=_optional_text(row.get("workdir")),
container=_optional_text(row.get("container")),
)
def _split_task_name(value: str) -> tuple[str, str | None]:
if value.endswith(")") and " (" in value:
process, _, tag = value.partition(" (")
return process, tag[:-1] or None
return value, None
def _normalize_state(value: Any) -> TaskState:
state = str(value or "unknown").strip().lower()
if state in {"queued", "submitted"}:
return "queued"
if state == "running":
return "running"
if state == "completed":
return "completed"
if state == "cached":
return "cached"
if state == "failed":
return "failed"
if state == "aborted":
return "aborted"
return "unknown"
def _optional_text(value: Any) -> str | None:
return str(value) if value is not None and value != "" else None
def _optional_int(value: Any) -> int | None:
try:
return int(value) if value is not None and value != "" else None
except (TypeError, ValueError):
return None
SHA-256: 56f37ed6e2c8c3597c2002d9b49c03b927b61ba90322fe2e0d3d1bb49d8ed752