Skip to content
Merged
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
1 change: 1 addition & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ Classmethods for construction:
- For project-level pipelines, `record_identifier` auto-defaults to `project_name` (which defaults to `"project"`)
- `force_overwrite` defaults to `True` at manager level; it is a settable property
- `result_formatter` is also a settable property (not an `__init__` param)
- `record_identifier` is a settable property for changing the default record after construction
- File/image results require `{"path": "...", "title": "..."}` dict format (`thumbnail_path` is optional for images)

## Schema-Free Mode
Expand Down
3 changes: 3 additions & 0 deletions pipestat/backends/db_backend/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
from .db_helpers import construct_db_url as construct_db_url
from .db_parsed_schema import ParsedSchemaDB as ParsedSchemaDB
from .dbbackend import DBBackend as DBBackend
26 changes: 19 additions & 7 deletions pipestat/backends/file_backend/filebackend.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,8 @@ def determine_results_file(self) -> None:
"""

if "{record_identifier}" in self.results_file_path:
# In the special case where the user wants to use {record_identifier} in file path
pass
# Template path: _data will be initialized per-record when resolved
self._data = None
else:
if not os.path.exists(self.results_file_path):
_LOGGER.debug(
Expand All @@ -105,6 +105,8 @@ def check_record_exists(
bool: Whether the record exists in the table.
"""

if self._data is None:
return False
return (
self.pipeline_name in self._data
and record_identifier in self._data[self.pipeline_name][self.pipeline_type]
Expand Down Expand Up @@ -148,8 +150,12 @@ def count_records(self) -> int:
Returns:
int: Number of records.
"""

return len(self._data[self.pipeline_name][self.pipeline_type])
if self._data is None:
return 0
try:
return len(self._data[self.pipeline_name][self.pipeline_type])
except (KeyError, TypeError):
return 0

def get_flag_file(
self, record_identifier: Optional[str] = None
Expand Down Expand Up @@ -193,9 +199,13 @@ def get_status(self, record_identifier: str) -> Optional[str]:
assert isinstance(flag_file, str), TypeError(
"Flag file path is expected to be a str, were multiple flags found?"
)
with open(flag_file, "r") as f:
status = f.read()
return status
try:
with open(flag_file, "r") as f:
status = f.read()
return status
except FileNotFoundError:
_LOGGER.debug(f"Flag file disappeared: {flag_file}")
return None
_LOGGER.debug(
f"Could not determine status for '{r_id}' record. "
f"No flags found in: {self.status_file_dir}"
Expand Down Expand Up @@ -239,6 +249,8 @@ def list_results(
"""
record_identifier = record_identifier or self.record_identifier

if self._data is None:
return []
try:
results = list(
self._data[self.pipeline_name][self.pipeline_type][record_identifier].keys()
Expand Down
1 change: 1 addition & 0 deletions pipestat/backends/pephub_backend/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
from .pephubbackend import PEPHUBBACKEND as PEPHUBBACKEND
1 change: 0 additions & 1 deletion pipestat/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,6 @@ class SchemaValidationErrorDuringReport(SchemaError):
"""Adds clarity to JSON schema validation errors by providing additional information to error message."""

def __init__(self, msg, record_identifier, result_identifier, result):

txt = msg # original schema validation error
txt += f"\nRecord identifier {record_identifier} \nResult_identifier {result_identifier} \nReported result: {result}"
super(SchemaValidationErrorDuringReport, self).__init__(txt)
Expand Down
54 changes: 30 additions & 24 deletions pipestat/pipestat.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,26 +51,7 @@
SchemaNotFoundError,
)
from .helpers import default_formatter, make_subdirectories, validate_type, zip_report
from .reports import HTMLReportBuilder, _create_stats_objs_summaries

try:
from pipestat.backends.db_backend.db_parsed_schema import ParsedSchemaDB as ParsedSchema
except ImportError:
from .parsed_schema import ParsedSchema

try:
from pipestat.backends.db_backend.db_helpers import construct_db_url
from pipestat.backends.db_backend.dbbackend import DBBackend
except ImportError:
# We let this pass, but if the user attempts to create DBBackend, check_dependencies raises exception.
pass

try:
from pipestat.backends.pephub_backend.pephubbackend import PEPHUBBACKEND
except ImportError:
# Let this pass, if phc dependencies cannot be imported, raise exception
pass

from .parsed_schema import ParsedSchema

_LOGGER = getLogger(PKG_NAME)

Expand Down Expand Up @@ -775,6 +756,14 @@ def initialize_pephubbackend(
record_identifier (str, optional): The record identifier.
pephub_path (str, optional): The path to the pephub registry.
"""
try:
from pipestat.backends.pephub_backend import PEPHUBBACKEND
except ImportError:
raise PipestatDependencyError(
msg="Missing required dependencies for PEPhub backend. "
"Install them with: pip install pipestat[pephub]"
)

self.backend = PEPHUBBACKEND(
record_identifier,
pephub_path,
Expand All @@ -785,10 +774,6 @@ def initialize_pephubbackend(
self.cfg[RESULT_FORMATTER],
)

@check_dependencies(
dependency_list=["DBBackend"],
msg="Missing required dependencies for this usage, e.g. try pip install pipestat['dbbackend']",
)
def initialize_dbbackend(
self, record_identifier: str | None = None, show_db_logs: bool = False
) -> None:
Expand All @@ -803,6 +788,14 @@ def initialize_dbbackend(
NoBackendSpecifiedError: If database configuration is missing.
PipestatDatabaseError: If database configuration is invalid.
"""
try:
from pipestat.backends.db_backend import DBBackend, ParsedSchemaDB, construct_db_url
except ImportError:
raise PipestatDependencyError(
msg="Missing required dependencies for database backend. "
"Install them with: pip install pipestat[dbbackend]"
)

_LOGGER.debug("Determined database as backend")
if not self.cfg.get(PROJECT_NAME):
raise ValueError(
Expand All @@ -827,6 +820,10 @@ def initialize_dbbackend(
raise PipestatDatabaseError(f"No database section ('{CFG_DATABASE_KEY}') in config")
self._show_db_logs = show_db_logs

# Re-parse schema with DB-aware parser for model building
if self._schema_path is not None:
self.cfg[SCHEMA_KEY] = ParsedSchemaDB(self._schema_path)

self.backend = DBBackend(
record_identifier,
self.cfg[PIPELINE_NAME],
Expand Down Expand Up @@ -1694,6 +1691,7 @@ def summarize(
Raises:
PipestatSummarizeError: If no results are found at the backend.
"""
from .reports import HTMLReportBuilder

if output_dir:
self.cfg[OUTPUT_DIR] = output_dir
Expand Down Expand Up @@ -1758,6 +1756,8 @@ def table(
Returns:
list[str]: File paths of the generated stats and objects files.
"""
from .reports import _create_stats_objs_summaries

if output_dir:
self.cfg[OUTPUT_DIR] = output_dir

Expand Down Expand Up @@ -1988,6 +1988,12 @@ def record_identifier(self) -> str | None:
"""
return self._resolve_record_identifier(None)

@record_identifier.setter
def record_identifier(self, value: str | None) -> None:
if value is not None and not value:
raise ValueError("record_identifier cannot be empty")
self.cfg[RECORD_IDENTIFIER] = value

@property
def record_count(self) -> int:
"""Number of records reported.
Expand Down
98 changes: 28 additions & 70 deletions pipestat/reports.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,9 @@
from logging import getLogger

import jinja2
import pandas as _pd
import yaml
from peppy.const import AMENDMENTS_KEY
from ubiquerg import mkabs
from ubiquerg import mkabs, parse_timedelta

from .const import (
BUTTON_APPEARANCE_BY_FLAG,
Expand Down Expand Up @@ -1273,43 +1272,6 @@ def _make_relpath(file_name, wd, context=None):
return relpath if not context else os.path.join(os.path.join(*context), relpath)


def _read_csv_encodings(path, encodings=["utf-8", "ascii"], **kwargs):
"""
Try to read file with the provided encodings.

Args:
path (str): Path to file.
encodings (list): List of encodings to try.
**kwargs: Additional keyword arguments.
"""
idx = 0
while idx < len(encodings):
e = encodings[idx]
try:
t = _pd.read_csv(path, encoding=e, **kwargs)
return t
except UnicodeDecodeError:
pass
idx = idx + 1
_LOGGER.warning(f"Could not read the log file '{path}' with encodings '{encodings}'")


def _read_tsv_to_json(path):
"""
Read a tsv file to a JSON formatted string.

Args:
path (str): To file path.

Returns:
str: JSON formatted string.
"""
assert os.path.exists(path), "The file '{}' does not exist".format(path)
_LOGGER.debug("Reading TSV from '{}'".format(path))
df = _pd.read_csv(path, sep="\t", index_col=False, header=None)
return df.to_json()


def fetch_pipeline_results(
project,
sample_name=None,
Expand Down Expand Up @@ -1453,10 +1415,21 @@ def _warn(what, e, sn):
0
] # Assumes the profile file will be in status dir
assert os.path.exists(profile), FileNotFoundError(f"Not found: {profile}")
df = _pd.read_csv(profile, sep="\t", comment="#", names=PROFILE_COLNAMES)
df["runtime"] = _pd.to_timedelta(df["runtime"])
times.append(_get_runtime(df))
mems.append(_get_maxmem(df))
rows = []
with open(profile) as fh:
reader = csv.reader(fh, delimiter="\t")
for row in reader:
line = row[0].strip() if row else ""
if not line or line.startswith("#"):
continue
if len(row) < len(PROFILE_COLNAMES):
continue
rows.append(dict(zip(PROFILE_COLNAMES, row)))
for r in rows:
r["runtime"] = parse_timedelta(r["runtime"])
r["mem"] = float(r["mem"])
times.append(_get_runtime(rows))
mems.append(_get_maxmem(rows))
except Exception as e:
_warn("profile", e, sample)
times.append(NO_DATA_PLACEHOLDER)
Expand Down Expand Up @@ -1496,35 +1469,20 @@ def create_glossary_table(project):
return render_jinja_template("glossary_table.html", get_jinja_env(), template_vars)


def _get_maxmem(profile: _pd.DataFrame) -> str:
"""
Get current peak memory.

Args:
profile (pandas.DataFrame): A data frame representing the current profile.tsv
for a sample.

Returns:
str: Max memory.
"""
return f"{str(max(profile['mem']) if not profile['mem'].empty else 0)} GB"

def _get_maxmem(rows: list[dict]) -> str:
"""Get peak memory across all profile rows."""
if not rows:
return "0 GB"
return f"{max(r['mem'] for r in rows)} GB"

def _get_runtime(profile_df: _pd.DataFrame) -> str:
"""
Collect the unique and last duplicated runtimes, sum them and then return in str format.

Args:
profile_df (pandas.DataFrame): A data frame representing the current profile.tsv
for a sample.

Returns:
str: Sum of runtimes.
"""
unique_df = profile_df[~profile_df.duplicated("cid", keep="last").values]
return str(
timedelta(seconds=sum(unique_df["runtime"].apply(lambda x: x.total_seconds())))
).split(".")[0]
def _get_runtime(rows: list[dict]) -> str:
"""Sum unique command runtimes, deduplicating reruns by cid (keeping last)."""
seen = {}
for r in rows:
seen[r["cid"]] = r
total = sum(r["runtime"].total_seconds() for r in seen.values())
return str(timedelta(seconds=total)).split(".")[0]


def get_file_for_table(prj, pipeline_name: str, appendix=None, directory=None) -> str:
Expand Down
7 changes: 3 additions & 4 deletions pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "pipestat"
version = "0.13.0"
version = "0.13.1"
description = "A pipeline results reporter"
readme = "README.md"
license = "BSD-2-Clause"
Expand All @@ -27,9 +27,8 @@ dependencies = [
"jsonschema",
"logmuse>=0.2.5",
"pyyaml",
"ubiquerg>=0.8.0",
"yacman>=0.9.5",
"pandas",
"ubiquerg>=0.9.1",
"yacman>=1.0.0",
"eido",
"jinja2",
]
Expand Down
22 changes: 22 additions & 0 deletions tests/test_multi_result_files.py
Original file line number Diff line number Diff line change
Expand Up @@ -83,3 +83,25 @@ def test_multi_results_summarize(
os.path.join(temp_dir, "aggregate_results.yaml")
)
assert r_id in data[psm.pipeline_name][psm.pipeline_type].keys()

@pytest.mark.parametrize("backend", ["file"])
def test_template_path_without_record_identifier(
self,
config_file_path,
results_file_path,
recursive_schema_file_path,
backend,
range_values,
):
"""PSM created with template path and no record_identifier can check status and list results."""
with TemporaryDirectory() as temp_dir:
results_file_path = os.path.join(temp_dir, "{record_identifier}/results.yaml")
psm = SamplePipestatManager(
results_file_path=results_file_path,
schema_path=recursive_schema_file_path,
)
# These are what looper calls on a PSM without record_identifier
assert psm.count_records() == 0
r_id = range_values[0][0]
assert psm.backend.check_record_exists(record_identifier=r_id) is False
assert psm.backend.list_results(record_identifier=r_id) == []
Loading