diff --git a/servicex_app/servicex_app/dataset_manager.py b/servicex_app/servicex_app/dataset_manager.py index 0db6b3463..284d368ff 100644 --- a/servicex_app/servicex_app/dataset_manager.py +++ b/servicex_app/servicex_app/dataset_manager.py @@ -123,6 +123,7 @@ def from_file_list( DatasetFile(paths=file, adler32="xxx", file_events=0, file_size=0) for file in file_list ] + dataset.n_files = len(file_list) logger.info( f"Upserted dataset for file list. Dataset Id is {dataset.id}", diff --git a/servicex_app/servicex_app/routes.py b/servicex_app/servicex_app/routes.py index 4b05fa15d..f419fb480 100644 --- a/servicex_app/servicex_app/routes.py +++ b/servicex_app/servicex_app/routes.py @@ -94,6 +94,8 @@ def add_routes( from servicex_app.web.transformation_request import transformation_request from servicex_app.web.transformation_results import transformation_results from servicex_app.web.multiple_codegen_list import multiple_codegen_list + from servicex_app.web.datasets import datasets as datasets_page + from servicex_app.web.dataset import dataset as dataset_page # Must be its own module to allow patching from servicex_app.web.create_profile import create_profile @@ -140,6 +142,8 @@ def add_routes( app.add_url_rule( "/multiple-codegen-list", "multiple_codegen_list", multiple_codegen_list ) + app.add_url_rule("/datasets", "datasets", datasets_page) + app.add_url_rule("/datasets/", "dataset", dataset_page) # User management and Authentication Endpoints api.add_resource(TokenRefresh, "/token/refresh") diff --git a/servicex_app/servicex_app/templates/base.html b/servicex_app/servicex_app/templates/base.html index e804d7985..799faac67 100644 --- a/servicex_app/servicex_app/templates/base.html +++ b/servicex_app/servicex_app/templates/base.html @@ -44,6 +44,11 @@ + {% if not config["ENABLE_AUTH"] or session['is_authenticated'] %} + + {% endif %} diff --git a/servicex_app/servicex_app/templates/dataset.html b/servicex_app/servicex_app/templates/dataset.html new file mode 100644 index 000000000..598b3ddf2 --- /dev/null +++ b/servicex_app/servicex_app/templates/dataset.html @@ -0,0 +1,135 @@ +{% extends "base.html" %} +{% block content %} +
+
+

Dataset

+
+ {% if not ds.stale %} + + {% elif ds.lookup_status and ds.lookup_status.value in ('bad_name', 'does_not_exist', 'internal_failure') %} + {{ ds.lookup_status.value }} + {% else %} + Flushed + {% endif %} +
+
+ +
+
ID
+
{{ ds.id }}
+ +
Name
+
+ {% if ds.did_finder == "user" %} + (file list) +
{{ ds.name }}
+ {% else %} + {{ ds.name }} + {% endif %} +
+ +
DID Finder
+
{{ ds.did_finder }}
+ +
Lookup Status
+
{{ ds.lookup_status.value if ds.lookup_status else '-' }}
+ +
Files
+
{{ humanize.intcomma(ds.n_files or 0) }}
+ +
Total Events
+
{{ humanize.intcomma(ds.events or 0) }}
+ +
Total Size
+
{{ humanize.naturalsize(ds.size or 0) }}
+ +
Last Used
+
+ {{ moment(ds.last_used).format("YYYY-MM-DD HH:mm:ss") if ds.last_used else "-" }} + +
+ +
Last Updated
+
+ {{ moment(ds.last_updated).format("YYYY-MM-DD HH:mm:ss") if ds.last_updated else "-" }} + +
+ +
Transform Requests
+
+ {% if ds.transform_requests %} + + {% else %} + - + {% endif %} +
+
+ +
Files ({{ humanize.intcomma(ds.files|length) }})
+ {% if ds.files %} +
+ + + + + + + + + + + + {% for f in ds.files %} + + + + + + + + {% endfor %} + +
IDPath(s)EventsSizeadler32
{{ f.id }}{{ f.paths }}{{ humanize.intcomma(f.file_events or 0) }}{{ humanize.naturalsize(f.file_size or 0) }}{{ f.adler32 or '-' }}
+
+ {% else %} +

No files.

+ {% endif %} +
+{% endblock %} + +{% block scripts %} + +{% endblock %} diff --git a/servicex_app/servicex_app/templates/dataset_table.html b/servicex_app/servicex_app/templates/dataset_table.html new file mode 100644 index 000000000..8894d64e6 --- /dev/null +++ b/servicex_app/servicex_app/templates/dataset_table.html @@ -0,0 +1,87 @@ +{% from 'bootstrap5/pagination.html' import render_pagination %} + +{% macro datasets_table(pagination, humanize) %} +
+ + + + + + + + + + + + + + + + + + {% for ds in pagination.items %} + + + + + + + + + + + + + {% endfor %} + +
All times in timezone: .
IDNameDID FinderLookup StatusFilesEventsSizeLast UsedLast UpdatedActions
{{ ds.id }} + + {% if ds.did_finder == "user" %}(file list){% else %}{{ ds.name }}{% endif %} + + {{ ds.did_finder }}{{ ds.lookup_status.value if ds.lookup_status else '-' }}{{ humanize.intcomma(ds.n_files or 0) }}{{ humanize.intcomma(ds.events or 0) }}{{ humanize.naturalsize(ds.size or 0) }}{{ moment(ds.last_used).format("YYYY-MM-DD HH:mm") if ds.last_used else "-" }}{{ moment(ds.last_updated).format("YYYY-MM-DD HH:mm") if ds.last_updated else "-" }} + {% if not ds.stale %} + + {% elif ds.lookup_status and ds.lookup_status.value in ('bad_name', 'does_not_exist', 'internal_failure') %} + {{ ds.lookup_status.value }} + {% else %} + Flushed + {% endif %} +
+ {% if pagination.items %} + {{ render_pagination(pagination, align='center') }} + {% else %} +
+ No datasets found. +
+ {% endif %} +
+{% endmacro %} + +{% macro datasets_table_scripts() %} + +{% endmacro %} diff --git a/servicex_app/servicex_app/templates/datasets.html b/servicex_app/servicex_app/templates/datasets.html new file mode 100644 index 000000000..e2b610795 --- /dev/null +++ b/servicex_app/servicex_app/templates/datasets.html @@ -0,0 +1,37 @@ +{% extends "base.html" %} + +{% from 'dataset_table.html' import datasets_table, datasets_table_scripts with context %} +{% from 'sort_dropdown.html' import sort_dropdown %} + +{% block content %} + +
+
+
+

Datasets

+
+
+ + +
+ + +
+
+ {{ sort_dropdown(dropdown_options, active_sort, active_order) }} +
+
+ {{ datasets_table(pagination, humanize) }} +
+
+ +{% endblock %} + +{% block scripts %} + {{ super() }} + {{ datasets_table_scripts() }} +{% endblock %} diff --git a/servicex_app/servicex_app/templates/sort_dropdown.html b/servicex_app/servicex_app/templates/sort_dropdown.html index e84ca4d05..ac5d9997d 100644 --- a/servicex_app/servicex_app/templates/sort_dropdown.html +++ b/servicex_app/servicex_app/templates/sort_dropdown.html @@ -2,11 +2,11 @@ diff --git a/servicex_app/servicex_app/web/dataset.py b/servicex_app/servicex_app/web/dataset.py new file mode 100644 index 000000000..030312466 --- /dev/null +++ b/servicex_app/servicex_app/web/dataset.py @@ -0,0 +1,12 @@ +from flask import render_template, abort + +from servicex_app.decorators import oauth_required +from servicex_app.models import Dataset + + +@oauth_required +def dataset(id_: int): + ds = Dataset.find_by_id(id_) + if not ds: + abort(404) + return render_template("dataset.html", ds=ds) diff --git a/servicex_app/servicex_app/web/datasets.py b/servicex_app/servicex_app/web/datasets.py new file mode 100644 index 000000000..893a5b1f4 --- /dev/null +++ b/servicex_app/servicex_app/web/datasets.py @@ -0,0 +1,63 @@ +import itertools + +from flask import render_template +from flask_restful import reqparse + +from servicex_app.decorators import oauth_required +from servicex_app.models import Dataset + +model_attributes = { + "last_used": Dataset.last_used, + "last_updated": Dataset.last_updated, + "name": Dataset.name, + "size": Dataset.size, + "events": Dataset.events, + "files": Dataset.n_files, +} +parser = reqparse.RequestParser() +parser.add_argument("page", default=1, type=int, location="args") +sort_choices = tuple(model_attributes.keys()) +parser.add_argument( + "sort", + choices=sort_choices, + default="last_used", + location="args", + help=f"Sort must be one of: {', '.join(map(repr, sort_choices))}.", +) +order_choices = ("asc", "desc") +parser.add_argument( + "order", + choices=order_choices, + default="desc", + location="args", + help="Order must be 'asc' or 'desc'.", +) +parser.add_argument( + "show_deleted", + type=lambda v: str(v).lower() in ("1", "true", "yes", "on"), + default=False, + location="args", +) + + +@oauth_required +def datasets(): + args = parser.parse_args() + sort, order = args["sort"], args["order"] + query = Dataset.query + if not args["show_deleted"]: + query = query.filter_by(stale=False) + + sort_column = model_attributes[sort] + sort_order = sort_column.asc() if order == "asc" else sort_column.desc() + pagination = query.order_by(sort_order).paginate( + page=args["page"], per_page=15, error_out=False + ) + return render_template( + "datasets.html", + pagination=pagination, + dropdown_options=list(itertools.product(sort_choices, order_choices)), + active_sort=sort, + active_order=order, + show_deleted=args["show_deleted"], + ) diff --git a/servicex_app/servicex_app_test/resources/internal/test_transform_file_complete.py b/servicex_app/servicex_app_test/resources/internal/test_transform_file_complete.py index 7de244df7..400424f32 100644 --- a/servicex_app/servicex_app_test/resources/internal/test_transform_file_complete.py +++ b/servicex_app/servicex_app_test/resources/internal/test_transform_file_complete.py @@ -29,7 +29,12 @@ import psycopg2 import pytest -from servicex_app.models import TransformationResult, TransformRequest, TransformStatus +from servicex_app.models import ( + TransformationResult, + TransformRequest, + TransformStatus, + DatasetFile, +) from servicex_app.transformer_manager import TransformerManager from servicex_app_test.resource_test_base import ResourceTestBase @@ -75,10 +80,20 @@ def othermock(self, mocker): return rv @pytest.fixture - def mock_transform_request_lookup(self, db_session, trqmock, othermock): + def dsfilemock(self, mocker): + rv = mocker.Mock() + rv.filter_by.return_value.with_for_update.return_value.one_or_none.return_value = ( # noqa: E501 + None + ) + return rv + + @pytest.fixture + def mock_transform_request_lookup(self, db_session, trqmock, othermock, dsfilemock): def switcher(cls): if cls == TransformRequest: return trqmock + elif cls == DatasetFile: + return dsfilemock else: return othermock diff --git a/servicex_app/servicex_app_test/web/test_dataset.py b/servicex_app/servicex_app_test/web/test_dataset.py new file mode 100644 index 000000000..889357015 --- /dev/null +++ b/servicex_app/servicex_app_test/web/test_dataset.py @@ -0,0 +1,56 @@ +from datetime import datetime, timezone + +from flask import Response, url_for +from pytest import fixture + +from servicex_app.models import Dataset, DatasetStatus + +from .web_test_base import WebTestBase + + +class TestDatasetDetail(WebTestBase): + + @staticmethod + def _fake_dataset(**overrides): + defaults = { + "id": 42, + "name": "rucio://data25/foo.bar", + "did_finder": "rucio", + "last_used": datetime.now(tz=timezone.utc), + "last_updated": datetime.now(tz=timezone.utc), + "n_files": 3, + "size": 1024, + "events": 100, + "lookup_status": DatasetStatus.complete, + "stale": False, + } + defaults.update(overrides) + ds = Dataset(**defaults) + ds.files = [] + ds.transform_requests = [] + return ds + + @fixture + def mock_find_by_id(self, mocker): + return mocker.patch("servicex_app.web.dataset.Dataset.find_by_id") + + def test_renders_existing_dataset( + self, client, user, mock_find_by_id, captured_templates + ): + ds = self._fake_dataset() + mock_find_by_id.return_value = ds + response: Response = client.get( + url_for("dataset", id_=42), headers=self.fake_header() + ) + assert response.status_code == 200 + mock_find_by_id.assert_called_once_with(42) + template, context = captured_templates[0] + assert template.name == "dataset.html" + assert context["ds"] is ds + + def test_missing_dataset_returns_404(self, client, user, mock_find_by_id): + mock_find_by_id.return_value = None + response: Response = client.get( + url_for("dataset", id_=999), headers=self.fake_header() + ) + assert response.status_code == 404 diff --git a/servicex_app/servicex_app_test/web/test_datasets.py b/servicex_app/servicex_app_test/web/test_datasets.py new file mode 100644 index 000000000..e330b3645 --- /dev/null +++ b/servicex_app/servicex_app_test/web/test_datasets.py @@ -0,0 +1,70 @@ +from flask import Response, url_for +from pytest import fixture + +from .web_test_base import WebTestBase + + +class TestDatasets(WebTestBase): + + @fixture + def mock_query(self, mocker): + mock_ds = mocker.patch("servicex_app.web.datasets.Dataset") + query = mock_ds.query + filtered = query.filter_by.return_value + return { + "raw": query.order_by.return_value, + "filtered": filtered.order_by.return_value, + } + + def test_default_filters_out_stale( + self, client, user, mock_query, captured_templates + ): + pagination = mock_query["filtered"].paginate( + page=1, per_page=15, total=0, items=[] + ) + mock_query["filtered"].paginate.return_value = pagination + response: Response = client.get(url_for("datasets"), headers=self.fake_header()) + assert response.status_code == 200 + template, context = captured_templates[0] + assert template.name == "datasets.html" + assert context["pagination"] == pagination + assert context["active_sort"] == "last_used" + assert context["active_order"] == "desc" + assert context["show_deleted"] is False + + def test_show_deleted_true_skips_stale_filter( + self, client, user, mock_query, captured_templates + ): + pagination = mock_query["raw"].paginate(page=1, per_page=15, total=0, items=[]) + mock_query["raw"].paginate.return_value = pagination + response: Response = client.get( + url_for("datasets") + "?show_deleted=true", + headers=self.fake_header(), + ) + assert response.status_code == 200 + template, context = captured_templates[0] + assert template.name == "datasets.html" + assert context["show_deleted"] is True + + def test_sort_and_order_are_applied( + self, client, user, mock_query, captured_templates + ): + pagination = mock_query["filtered"].paginate( + page=1, per_page=15, total=0, items=[] + ) + mock_query["filtered"].paginate.return_value = pagination + response: Response = client.get( + url_for("datasets") + "?sort=events&order=asc", + headers=self.fake_header(), + ) + assert response.status_code == 200 + template, context = captured_templates[0] + assert context["active_sort"] == "events" + assert context["active_order"] == "asc" + + def test_invalid_sort_choice_returns_400(self, client, user, mock_query): + response: Response = client.get( + url_for("datasets") + "?sort=bogus", + headers=self.fake_header(), + ) + assert response.status_code == 400 diff --git a/transformer_sidecar/src/transformer_sidecar/transformer_stats/aod_stats.py b/transformer_sidecar/src/transformer_sidecar/transformer_stats/aod_stats.py index eec42a828..9f53950eb 100644 --- a/transformer_sidecar/src/transformer_sidecar/transformer_stats/aod_stats.py +++ b/transformer_sidecar/src/transformer_sidecar/transformer_stats/aod_stats.py @@ -35,9 +35,15 @@ class AODStats(TransformerStats): def __init__(self, log_path: Path): super().__init__(log_path) - matches = re.findall(r"Processed (\d+) events", self.log_body) - if len(matches) == 1: - self.total_events = int(matches[0]) + # First see if the total events processed is printed + range_matches = re.findall(r"Processing events 0-(\d+) in file", self.log_body) + if range_matches: + self.total_events = int(range_matches[-1]) + else: + # Fallback: "Processed N events" is a periodic progress marker + progress_matches = re.findall(r"Processed (\d+) events", self.log_body) + if progress_matches: + self.total_events = max(int(m) for m in progress_matches) # Look for incorrect property names matches = re.findall( diff --git a/transformer_sidecar/tests/transformer_stats/test_aod_stats.py b/transformer_sidecar/tests/transformer_stats/test_aod_stats.py index 06dbaa6d7..7ef68b60b 100644 --- a/transformer_sidecar/tests/transformer_stats/test_aod_stats.py +++ b/transformer_sidecar/tests/transformer_stats/test_aod_stats.py @@ -52,6 +52,68 @@ def test_aod_stats(): os.remove(test_logfile_path) +def test_aod_stats_large_file(): + with tempfile.NamedTemporaryFile(mode="w", delete=False) as fp: + test_logfile_path = Path(fp.name) + fp.write( + "Package.EventLoop INFO created submission directory " + "/home/atlas/rel/build/bogus\n" + "Package.EventLoop INFO submitting job in " + "/home/atlas/rel/build/bogus\n" + "Package.EventLoop INFO Running sample: ANALYSIS\n" + "Package.EventLoop INFO xAODInput = 1\n" + "Package.EventLoop INFO calling firstInitialize on all modules\n" + "Package.EventLoop INFO calling preFileInitialize on all modules\n" + "Package.EventLoop INFO Opening file " + "root://example/DAOD_PHYSLITE.pool.root.1\n" + "Package.EventLoop INFO Processing events 0-213934 in file " + "root://example/DAOD_PHYSLITE.pool.root.1\n" + "Package.EventLoop INFO Processed 10000 events\n" + "Package.EventLoop INFO Processed 20000 events\n" + "Package.EventLoop INFO Processed 100000 events\n" + "Package.EventLoop INFO Processed 200000 events\n" + "Package.EventLoop INFO Processed 210000 events\n" + "LeakCheckModule INFO Memory increase/change during the job:\n" + "Package.EventLoop INFO worker finished successfully\n" + "Package.EventLoop INFO done\n" + ) + fp.close() + aod_stats = AODStats(test_logfile_path) + assert aod_stats.total_events == 213934 + os.remove(test_logfile_path) + + +def test_aod_stats_small_file(): + with tempfile.NamedTemporaryFile(mode="w", delete=False) as fp: + test_logfile_path = Path(fp.name) + fp.write( + "Package.EventLoop INFO Opening file " + "root://example/DAOD_PHYSLITE.pool.root.1\n" + "Package.EventLoop INFO Processing events 0-11860 in file " + "root://example/DAOD_PHYSLITE.pool.root.1\n" + "Package.EventLoop INFO Processed 10000 events\n" + "Package.EventLoop INFO worker finished successfully\n" + ) + fp.close() + aod_stats = AODStats(test_logfile_path) + assert aod_stats.total_events == 11860 + os.remove(test_logfile_path) + + +def test_aod_stats_fallback_progress_markers_only(): + with tempfile.NamedTemporaryFile(mode="w", delete=False) as fp: + test_logfile_path = Path(fp.name) + fp.write( + "Package.EventLoop INFO Processed 10000 events\n" + "Package.EventLoop INFO Processed 20000 events\n" + "Package.EventLoop INFO Processed 30000 events\n" + ) + fp.close() + aod_stats = AODStats(test_logfile_path) + assert aod_stats.total_events == 30000 + os.remove(test_logfile_path) + + def test_bad_property(): with tempfile.NamedTemporaryFile(mode="w", delete=False) as fp: test_logfile_path = Path(fp.name)