|
| 1 | +import json |
| 2 | +import threading |
| 3 | +from pathlib import Path |
| 4 | + |
| 5 | +import yaml |
| 6 | + |
| 7 | +from ... import util |
| 8 | +from .. import Aggregator, AggregatorError |
| 9 | +from ..jsonl.jsonl import verbatim_move |
| 10 | + |
| 11 | +get_logger = util.get_loggers("atex.aggregator.yamld") |
| 12 | + |
| 13 | + |
| 14 | +class YamlDocumentAggregator(Aggregator): |
| 15 | + """ |
| 16 | + - `target` is a string/Path to a `.yaml` file for all ingested results |
| 17 | + to be aggregated (written) to. |
| 18 | +
|
| 19 | + - `files` is a string/Path of the top-level parent for all per-platform |
| 20 | + / per-test files uploaded by tests. |
| 21 | +
|
| 22 | + - `allow_duplicate` permits any one test name to be ingested more than |
| 23 | + once, appending ` (1)` to the second test name entry, ` (2)` to the |
| 24 | + third, etc. |
| 25 | + """ |
| 26 | + |
| 27 | + def __init__(self, target, files, *, allow_duplicate=False): |
| 28 | + self.lock = threading.RLock() |
| 29 | + self.logger = get_logger() |
| 30 | + |
| 31 | + self.target = Path(target) |
| 32 | + self.files = Path(files) |
| 33 | + self.allow_duplicate = allow_duplicate |
| 34 | + self.seen_tests = {} |
| 35 | + self.target_fobj = None |
| 36 | + |
| 37 | + def start(self): |
| 38 | + self.logger.debug(f"starting: {self}") |
| 39 | + |
| 40 | + if self.target.exists(follow_symlinks=False): |
| 41 | + raise FileExistsError(f"{self.target} already exists") |
| 42 | + self.target_fobj = open(self.target, "w") |
| 43 | + |
| 44 | + if self.files.exists(follow_symlinks=False): |
| 45 | + raise FileExistsError(f"{self.files} already exists") |
| 46 | + self.files.mkdir() |
| 47 | + |
| 48 | + def stop(self): |
| 49 | + self.logger.debug(f"stopping: {self}") |
| 50 | + |
| 51 | + if self.target_fobj: |
| 52 | + self.target_fobj.close() |
| 53 | + self.target_fobj = None |
| 54 | + |
| 55 | + def ingest(self, platform, test_name, artifacts): |
| 56 | + unique_id = (platform, test_name) |
| 57 | + with self.lock: |
| 58 | + if unique_id in self.seen_tests: |
| 59 | + if not self.allow_duplicate: |
| 60 | + raise AggregatorError( |
| 61 | + f"'{test_name}' was already ingested once for '{platform}'", |
| 62 | + ) |
| 63 | + else: |
| 64 | + test_name = f"{test_name} ({self.seen_tests[unique_id]})" |
| 65 | + self.seen_tests[unique_id] += 1 |
| 66 | + else: |
| 67 | + self.seen_tests[unique_id] = 1 |
| 68 | + |
| 69 | + self.logger.info(f"ingesting '{platform}' / '{test_name}' from '{artifacts}'") |
| 70 | + |
| 71 | + artifacts = Path(artifacts) |
| 72 | + artifacts_results = artifacts / "results" |
| 73 | + artifacts_files = artifacts / "files" |
| 74 | + |
| 75 | + if not artifacts_results.exists(follow_symlinks=False): |
| 76 | + raise FileNotFoundError(f"{artifacts_results} does not exist") |
| 77 | + |
| 78 | + platform_files = self.files / util.normalize_path(platform) |
| 79 | + target_test_files = platform_files / util.normalize_path(test_name) |
| 80 | + if target_test_files.exists(follow_symlinks=False): |
| 81 | + raise FileExistsError(f"{target_test_files} already exists for {test_name}") |
| 82 | + |
| 83 | + # any None or empty values are deleted later, |
| 84 | + # to preserve dict insertion order with these on top |
| 85 | + document = { |
| 86 | + "platform": platform, |
| 87 | + "name": test_name, |
| 88 | + "status": None, |
| 89 | + "note": None, |
| 90 | + "files": [], |
| 91 | + "subtests": [], |
| 92 | + } |
| 93 | + |
| 94 | + with open(artifacts_results) as f: |
| 95 | + for raw_line in f: |
| 96 | + result_line = json.loads(raw_line) |
| 97 | + |
| 98 | + # if it is a subtest, add it to subtests |
| 99 | + if name := result_line.get("name"): |
| 100 | + subtest = {"name": name} |
| 101 | + if status := result_line.get("status"): |
| 102 | + subtest["status"] = status |
| 103 | + if files := result_line.get("files"): |
| 104 | + subtest["files"] = files |
| 105 | + if note := result_line.get("note"): |
| 106 | + subtest["note"] = note |
| 107 | + document["subtests"].append(subtest) |
| 108 | + |
| 109 | + # update document for the test itself |
| 110 | + else: |
| 111 | + if status := result_line.get("status"): |
| 112 | + document["status"] = status |
| 113 | + if note := result_line.get("note"): |
| 114 | + document["note"] = note |
| 115 | + |
| 116 | + file_names = [] |
| 117 | + # process the file specified by the 'testout' key |
| 118 | + if "testout" in result_line: |
| 119 | + file_names.append(result_line["testout"]) |
| 120 | + # process any additional files in the 'files' key |
| 121 | + if files := result_line.get("files"): |
| 122 | + file_names += files |
| 123 | + if file_names: |
| 124 | + document["files"] += file_names |
| 125 | + |
| 126 | + if document["status"] is None: |
| 127 | + del document["status"] |
| 128 | + if document["note"] is None: |
| 129 | + del document["note"] |
| 130 | + if not document["files"]: |
| 131 | + del document["files"] |
| 132 | + if not document["subtests"]: |
| 133 | + del document["subtests"] |
| 134 | + |
| 135 | + with self.lock: |
| 136 | + yaml.dump( |
| 137 | + document, |
| 138 | + self.target_fobj, |
| 139 | + explicit_start=True, |
| 140 | + default_flow_style=False, |
| 141 | + sort_keys=False, |
| 142 | + ) |
| 143 | + self.target_fobj.flush() |
| 144 | + |
| 145 | + # clean up the source test_results (Aggregator should 'mv', not 'cp') |
| 146 | + Path(artifacts_results).unlink() |
| 147 | + |
| 148 | + # if the artifacts files directory is not empty |
| 149 | + if any(artifacts_files.iterdir()): |
| 150 | + platform_files.mkdir(exist_ok=True) |
| 151 | + # TODO: why does this work without .mkdir(target_test_files.parent) ? |
| 152 | + verbatim_move(artifacts_files, target_test_files) |
| 153 | + |
| 154 | + def __str__(self): |
| 155 | + class_name = self.__class__.__name__ |
| 156 | + return f"{class_name}({str(self.target)}, {str(self.files)})" |
0 commit comments