Skip to content

Run QA Workflow esgf_qa.run_qa

Public entry point and high-level orchestration for esgqa.

CheckerMetadata dataclass

Resolved version and package provenance for one checker.

Source code in esgf_qa/checker_registry.py
13
14
15
16
17
18
19
@dataclass(frozen=True)
class CheckerMetadata:
    """Resolved version and package provenance for one checker."""

    checker_version: str
    package_name: str | None = None
    package_version: str | None = None

FileInventory dataclass

Files and mappings consumed by the QA workflow.

Source code in esgf_qa/discovery.py
17
18
19
20
21
22
23
24
25
26
27
28
@dataclass
class FileInventory:
    """Files and mappings consumed by the QA workflow."""

    files: list[str]
    file_details: dict
    dataset_files: dict
    directory_datasets: dict
    checker_options: dict
    discovered_file_count: int = 0
    blacklisted_files: dict[str, list[str]] = field(default_factory=dict)
    not_whitelisted_files: list[str] = field(default_factory=list)

RunConfig dataclass

Validated configuration required by the QA workflow.

Source code in esgf_qa/cli.py
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
@dataclass
class RunConfig:
    """Validated configuration required by the QA workflow."""

    parent_dir: str
    result_dir: str
    checkers: list[str]
    info: str
    resume: bool
    include_consistency_checks: bool
    checker_options: dict
    parallel_processes: int
    time_checks_only: bool
    progress_file: Path
    dataset_file: Path
    processed_files: set[str]
    processed_datasets: set[str]
    whitelist: list[str] = field(default_factory=list)
    blacklist: list[str] = field(default_factory=list)
    rerun_all: bool = False

call_process_dataset(args)

Unpack multiprocessing arguments for :func:process_dataset.

Source code in esgf_qa/workers.py
355
356
357
def call_process_dataset(args):
    """Unpack multiprocessing arguments for :func:`process_dataset`."""
    return process_dataset(*args)

call_process_file(args)

Unpack multiprocessing arguments for :func:process_file.

Source code in esgf_qa/workers.py
350
351
352
def call_process_file(args):
    """Unpack multiprocessing arguments for :func:`process_file`."""
    return process_file(*args)

format_checker_version(checker, metadata)

Format a checker specification with its providing package version.

Source code in esgf_qa/checker_registry.py
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
def format_checker_version(checker, metadata):
    """Format a checker specification with its providing package version."""
    checker_name = checker.split(":", 1)[0]
    checker_metadata = metadata[checker_name]
    checker_label = (
        f"{checker_dict.get(checker_name, '')} "
        f"{checker_name}:{checker_metadata.checker_version}"
    ).strip()
    if checker_metadata.package_name is not None:
        checker_label += (
            f" ({checker_metadata.package_name} " f"{checker_metadata.package_version})"
        )
    return checker_label

get_checker_metadata(checkers, checker_options=None)

Resolve checker versions and their providing distributions.

Source code in esgf_qa/checker_registry.py
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
def get_checker_metadata(checkers, checker_options=None):
    """Resolve checker versions and their providing distributions."""
    check_suite = CheckSuite(options=checker_options or {})
    check_suite.load_all_available_checkers()

    checker_packages = {}
    for checker_entry_point in entry_points(group="compliance_checker.suites"):
        try:
            checker_obj = checker_entry_point.load()
        except Exception:
            continue
        checker_name = getattr(checker_obj, "_cc_spec", None) or getattr(
            checker_obj, "name", None
        )
        checker_version = str(getattr(checker_obj, "_cc_spec_version", "unknown"))
        distribution = getattr(checker_entry_point, "dist", None)
        if checker_name and distribution is not None:
            checker_packages[(checker_name, checker_version)] = (
                distribution.name,
                distribution.version,
            )

    metadata = {}
    for checker in checkers:
        checker_name = checker.split(":", 1)[0]
        if checker_name in checker_dict_ext and checker_name not in checker_dict:
            metadata[checker_name] = CheckerMetadata(
                esgf_qa_version, "esgf-qa", esgf_qa_version
            )
            continue

        checker_obj = check_suite.checkers.get(checker)
        if checker_obj is None:
            prefix = checker_name + ":"
            candidates = [key for key in check_suite.checkers if key.startswith(prefix)]
            if candidates:
                resolved_key = max(
                    candidates,
                    key=lambda key: pversion.parse(key.split(":", 1)[1]),
                )
                checker_obj = check_suite.checkers.get(resolved_key)
        checker_version = (
            str(checker_obj._cc_spec_version) if checker_obj is not None else "unknown"
        )
        package = checker_packages.get((checker_name, checker_version))
        metadata[checker_name] = CheckerMetadata(
            checker_version,
            package[0] if package else None,
            package[1] if package else None,
        )
    return metadata

get_default_result_dir()

Return the timestamped default result directory.

Source code in esgf_qa/run_qa.py
53
54
55
56
def get_default_result_dir():
    """Return the timestamped default result directory."""
    result_hash = hashlib.md5(_timestamp_with_ms.encode()).hexdigest()
    return os.path.abspath(f"esgf-qa-results_{_timestamp_filename}_{result_hash}")

get_dsid(files_to_check_dict, dataset_files_map_ext, file_path, project_ids)

Build a dataset identifier from its directory and filename.

Source code in esgf_qa/discovery.py
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
def get_dsid(files_to_check_dict, dataset_files_map_ext, file_path, project_ids):
    """Build a dataset identifier from its directory and filename."""
    directory_parts = files_to_check_dict[file_path]["id_dir"].split("/")
    filename_parts = files_to_check_dict[file_path]["id_fn"].split("_")
    dataset_id = ".".join(directory_parts)
    lower_parts = [part.lower() for part in directory_parts]
    for project_id in project_ids:
        if project_id in lower_parts:
            last_index = len(lower_parts) - 1 - lower_parts[::-1].index(project_id)
            dataset_id = ".".join(directory_parts[last_index:])
            break
    directory = files_to_check_dict[file_path]["id_dir"]
    if len(dataset_files_map_ext[directory]) > 1:
        dataset_id += "." + ".".join(filename_parts)
    return dataset_id

get_installed_checker_versions()

Return installed checker versions grouped by checker name.

Source code in esgf_qa/checker_registry.py
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
def get_installed_checker_versions():
    """Return installed checker versions grouped by checker name."""
    check_suite = CheckSuite()
    check_suite.load_all_available_checkers()
    installed_versions = {}
    for checker in check_suite.checkers:
        try:
            name, checker_version = checker.split(":")
        except ValueError:
            name, checker_version = checker, "latest"
        if checker_version == "latest":
            continue
        installed_versions.setdefault(name, []).append(checker_version)
    for name, versions in installed_versions.items():
        installed_versions[name] = sorted(versions, key=pversion.parse) + ["latest"]
    return installed_versions

main(argv=None)

Run the complete QA command-line workflow.

Source code in esgf_qa/run_qa.py
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
def main(argv=None):
    """Run the complete QA command-line workflow."""
    config = prepare_run(get_default_result_dir(), argv)
    inventory = discover_files(config)
    excluded_report = write_excluded_files(inventory, config)
    reconcile_resume_inventory(inventory, config)
    if not inventory.files:
        if inventory.discovered_file_count == 0:
            raise FileNotFoundError(
                f"No NetCDF files found under '{config.parent_dir}'."
            )
        raise RuntimeError(
            "No files remain to check: "
            f"{inventory.discovered_file_count} NetCDF files were discovered, "
            f"{format_exclusion_counts(inventory, config)}. "
            f"See '{excluded_report}'."
        )
    write_inventory(inventory, config.result_dir)
    summary, reference_datasets = run_workflow(config, inventory)
    write_results(
        config,
        inventory,
        summary,
        reference_datasets,
        _timestamp_with_ms,
        _timestamp_filename,
        _timestamp_pprint,
    )

normalize_checker_specs(checkers_versions)

Normalize latest versions to Compliance Checker's unversioned form.

Source code in esgf_qa/checker_registry.py
108
109
110
111
112
113
def normalize_checker_specs(checkers_versions):
    """Normalize latest versions to Compliance Checker's unversioned form."""
    return sorted(
        checker if checker_version == "latest" else f"{checker}:{checker_version}"
        for checker, checker_version in checkers_versions.items()
    )

parse_options(opts)

Parse checker:option[:value] CLI options into a nested mapping.

Source code in esgf_qa/cli.py
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
def parse_options(opts):
    """Parse ``checker:option[:value]`` CLI options into a nested mapping."""
    options_dict = defaultdict(dict)
    for option in opts:
        parts = option.split(":", 2)
        if len(parts) < 2 or not parts[0] or not parts[1]:
            raise ValueError(
                f"Could not split option '{option}', seems illegally formatted. "
                "The required format is '<checker>:<option_name>[:<option_value>]', "
                "for example 'mip:tables:/path/to/Tables'."
            )
        checker_type, checker_option, *checker_value = parts
        checker_value = checker_value[0] if checker_value else True
        options_dict[checker_type][checker_option] = checker_value
    return options_dict

process_dataset(dataset_id, dataset_files, checkers, checker_options, file_details, was_processed)

Run or reuse dataset-level consistency checks.

Source code in esgf_qa/workers.py
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
def process_dataset(
    dataset_id,
    dataset_files,
    checkers,
    checker_options,
    file_details,
    was_processed,
):
    """Run or reuse dataset-level consistency checks."""
    dataset_files_map = {dataset_id: dataset_files}
    result_file = file_details[dataset_files[0]]["result_file_ds"]
    result = get_reusable_dataset_result(checkers, result_file, was_processed)
    if result is not None:
        print(f"Read result from disk for '{dataset_id}'.")
        return dataset_id, result
    if was_processed:
        print(f"Rerunning previously erroneous checks for '{dataset_id}'.")
    else:
        print(f"Running checks for '{dataset_id}'.")

    _remove_stale_output(result_file)
    result = {}
    for checker_spec in checkers:
        checker = checker_spec.split(":", 1)[0]
        checker_fct = DATASET_CHECKERS.get(checker)
        if checker_fct is None:
            result[checker] = {
                "errors": {
                    checker: {
                        "msg": f"Checker '{checker}' not found.",
                        "files": dataset_files,
                    }
                }
            }
            continue
        try:
            result[checker] = checker_fct(
                dataset_id,
                dataset_files_map,
                file_details,
                checker_options[checker],
            )
        except Exception as error:
            result[checker] = _dataset_check_runtime_error(
                checker_fct.__name__, error, dataset_files
            )

    _replace_json(result_file, result)
    return dataset_id, result

process_file(file_path, checkers, checker_options, file_details, was_processed)

Run or reuse file-level checks for one file.

Source code in esgf_qa/workers.py
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
def process_file(
    file_path,
    checkers,
    checker_options,
    file_details,
    was_processed,
):
    """Run or reuse file-level checks for one file."""
    consistency_file = file_details["consistency_file"]
    result_file = file_details["result_file"]
    result = get_reusable_file_result(checkers, file_details, was_processed)
    if result is not None:
        print(f"Read result from disk for '{file_path}'.")
        return file_path, result
    if was_processed:
        print(f"Rerunning incomplete or previously erroneous checks for '{file_path}'.")
    else:
        print(f"Running checks for '{file_path}'.")

    # A rerun must not mistake output from an earlier attempt for freshly
    # generated output when the checker fails before writing its new files.
    _remove_stale_output(result_file)
    _remove_stale_output(consistency_file)
    result = run_compliance_checker(file_path, checkers, checker_options)
    check_results = {}
    for checker_spec in checkers:
        checker = checker_spec.split(":", 1)[0]
        checker_result = {"errors": {}}
        check_results[checker] = checker_result
        for check in result[checker_spec][0]:
            serialized_check = {
                "weight": check.weight,
                "value": check.value,
                "msgs": check.msgs,
                "method": check.check_method,
                "children": check.children,
            }
            previous_result = checker_result.get(check.name)
            if previous_result is None:
                checker_result[check.name] = serialized_check
            elif isinstance(previous_result, list):
                previous_result.append(serialized_check)
            else:
                # A single Compliance Checker method can report the same named
                # check at multiple severities; retain every result record.
                checker_result[check.name] = [previous_result, serialized_check]

        # Error keys are checker method names, whereas normal result keys above
        # are the human-readable check names exposed by Compliance Checker.
        for check_method, error_details in result[checker_spec][1].items():
            checker_result["errors"][check_method] = (
                _format_compliance_checker_runtime_error(check_method, error_details)
            )

    if not os.path.isfile(consistency_file):
        for checker in (
            checker.split(":", 1)[0]
            for checker in checkers
            if checker.split(":", 1)[0] in checker_supporting_consistency_checks
        ):
            check_results[checker]["errors"][
                "consistency_output"
            ] = f"Expected consistency output file was not created: '{consistency_file}'."

    _replace_json(result_file, check_results)
    return file_path, check_results

reconcile_resume_inventory(inventory, config)

Report resume inventory changes and invalidate results for missing files.

Source code in esgf_qa/resume.py
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
def reconcile_resume_inventory(inventory, config):
    """Report resume inventory changes and invalidate results for missing files."""
    if not config.resume:
        return None

    previous_inventory_path = Path(config.result_dir, "files_to_check.json")
    previous_details_path = Path(config.result_dir, "files_to_check_dict.json")
    report_path = Path(config.result_dir, "resume_inventory_changes.json")
    try:
        with open(previous_inventory_path) as file:
            previous_files = json.load(file)
    except (OSError, json.JSONDecodeError) as error:
        warnings.warn(
            "Could not compare the resumed file inventory because "
            f"'{previous_inventory_path}' could not be read: {error}"
        )
        return None
    if not isinstance(previous_files, list) or not all(
        isinstance(file_path, str) for file_path in previous_files
    ):
        warnings.warn(
            "Could not compare the resumed file inventory because "
            f"'{previous_inventory_path}' is not a list of file paths."
        )
        return None

    previous_file_set = set(previous_files)
    current_file_set = set(inventory.files)
    added_files = sorted(current_file_set - previous_file_set)
    missing_files = sorted(previous_file_set - current_file_set)
    affected_datasets = set()
    if missing_files:
        try:
            with open(previous_details_path) as file:
                previous_details = json.load(file)
            if not isinstance(previous_details, dict):
                raise TypeError("expected a dictionary")
            for file_path in missing_files:
                details = previous_details.get(file_path)
                if not isinstance(details, dict) or not isinstance(
                    details.get("id"), str
                ):
                    raise TypeError(f"missing dataset identifier for '{file_path}'")
                affected_datasets.add(details["id"])
        except (OSError, json.JSONDecodeError, KeyError, TypeError) as error:
            # Without the previous file-to-dataset mapping, no dataset result can
            # safely be assumed independent of the files that disappeared.
            affected_datasets = set(config.processed_datasets)
            warnings.warn(
                "Could not determine which cached dataset results contain files "
                f"that are no longer found; invalidating all dataset results: {error}"
            )

        config.processed_files.difference_update(missing_files)
        config.processed_datasets.difference_update(affected_datasets)
        write_progress(config.progress_file, config.processed_files)
        write_progress(config.dataset_file, config.processed_datasets)

    report = {
        "summary": {
            "previous_selected": len(previous_file_set),
            "current_selected": len(current_file_set),
            "added": len(added_files),
            "no_longer_found": len(missing_files),
        },
        "added": added_files,
        "no_longer_found": missing_files,
    }
    with open(report_path, "w") as file:
        json.dump(report, file, indent=4)

    if added_files or missing_files:
        print("Resume inventory changed:")
        if added_files:
            label = "file" if len(added_files) == 1 else "files"
            print(f" - {len(added_files)} new {label} will be checked.")
        if missing_files:
            label = "file is" if len(missing_files) == 1 else "files are"
            print(
                f" - {len(missing_files)} previously selected {label} no longer found."
            )
    else:
        print("Resume inventory is unchanged.")
    print(f"Resume inventory comparison was saved to '{report_path}'.")
    return str(report_path)

run_compliance_checker(file_path, checkers, checker_options=None)

Run Compliance Checker for one file, isolating checker-level failures.

Source code in esgf_qa/workers.py
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
def run_compliance_checker(file_path, checkers, checker_options=None):
    """Run Compliance Checker for one file, isolating checker-level failures."""
    checker_options = checker_options or {}
    try:
        check_suite = CheckSuite(options=checker_options)
        _ensure_checkers_loaded(check_suite)
    except Exception as error:
        return {
            checker: _failed_checker_result(error, "run_compliance_checker")
            for checker in checkers
        }

    try:
        dataset = check_suite.load_dataset(file_path)
    except Exception as error:
        return {
            checker: _failed_checker_result(error, "load_dataset")
            for checker in checkers
        }

    time_checks_only = checker_options.get("cc6", {}).get(
        "time_checks_only", False
    ) or checker_options.get("mip", {}).get("time_checks_only", False)
    include_checks = (
        ["check_time_continuity", "check_time_bounds", "check_time_range"]
        if time_checks_only
        else None
    )

    results = {}
    close_error = None
    try:
        for checker in checkers:
            checker_include = (
                include_checks if checker.split(":", 1)[0] in {"cc6", "mip"} else None
            )
            try:
                checker_results = check_suite.run_all(
                    dataset,
                    [checker],
                    include_checks=checker_include,
                    skip_checks=[],
                )
                if checker not in checker_results:
                    raise RuntimeError(
                        f"Compliance Checker returned no result for '{checker}'."
                    )
                results[checker] = checker_results[checker]
            except Exception as error:
                results[checker] = _failed_checker_result(error, "run_checker")
    finally:
        if hasattr(dataset, "close"):
            try:
                dataset.close()
            except Exception as error:
                close_error = error

    if close_error is not None:
        for checker in checkers:
            results.setdefault(
                checker, _failed_checker_result(close_error, "close_dataset")
            )[1]["close_dataset"] = (close_error, close_error.__traceback__)
    return results

run_dataset_collection_check(summary, checker, checker_fct, ds_map, files_to_check_dict, checker_options)

Run an all-dataset check and aggregate a runtime error if it fails.

Source code in esgf_qa/workers.py
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
def run_dataset_collection_check(
    summary,
    checker,
    checker_fct,
    ds_map,
    files_to_check_dict,
    checker_options,
):
    """Run an all-dataset check and aggregate a runtime error if it fails."""
    try:
        return checker_fct(ds_map, files_to_check_dict, checker_options)
    except Exception as error:
        for dataset_id, files in ds_map.items():
            summary.update_ds(
                {
                    checker: _dataset_check_runtime_error(
                        checker_fct.__name__, error, files
                    )
                },
                dataset_id,
            )
        return None

track_checked_datasets(checked_datasets_file, checked_datasets)

Append dataset identifiers to a progress file.

Source code in esgf_qa/resume.py
345
346
347
348
349
350
def track_checked_datasets(checked_datasets_file, checked_datasets):
    """Append dataset identifiers to a progress file."""
    with open(checked_datasets_file, "a") as file:
        writer = csv.writer(file)
        for dataset_id in checked_datasets:
            writer.writerow([dataset_id])