Skip to content

Commit d8211d8

Browse files
committed
uploader: clarify runtime data flow
Distinguish HTTP outcomes from file results, give resource and telemetry state domain-specific names, and replace string-based report aggregation with explicit fields. Remove stale rule commentary while preserving runtime and report contracts. Validated with the 449-test Python suite, all 199 //tools tests, compileall, template lint, buildifier, and git diff checks.
1 parent de4b02e commit d8211d8

6 files changed

Lines changed: 93 additions & 97 deletions

File tree

tools/core/test_optimization_uploader.bzl

Lines changed: 2 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -32,22 +32,10 @@ Why generated launchers/configuration:
3232
Developer navigation:
3333
- Starlark test helpers (manifest/CODEOWNERS parsing parity) near the top.
3434
- `_uploader_impl` builds both script templates and wires rule outputs.
35-
- Script templates contain the runtime behavior for discovery, enrichment,
36-
uploads, retries, and locking.
35+
- `uploader_py` owns current runtime behavior; legacy generated scripts remain
36+
here only as a rollout fallback.
3737
"""
3838

39-
# Usage pattern:
40-
# bazel test //... || test_status=$?; test_status=${test_status:-0}; DD_API_KEY="$DD_API_KEY" DD_SITE="$DD_SITE" bazel run //:dd_upload_payloads; exit $test_status
41-
#
42-
# Key features:
43-
# - Discovers all test.outputs/ directories in bazel-testlogs automatically
44-
# - Supports sharded tests (shard_N_of_M/) and retries (run_N_of_M/)
45-
# - Uploads test payloads to CI Test Cycle intake
46-
# - Uploads coverage payloads to Code Coverage intake
47-
# - Deletes payloads after successful upload (unless DD_TEST_OPTIMIZATION_KEEP_PAYLOADS=1)
48-
# - Uses workspace-level lock to prevent concurrent uploaders
49-
# - Enriches payloads with context.json metadata
50-
5139
load(
5240
"//tools/core:common_utils.bzl",
5341
"RULES_VERSION",

tools/core/uploader_py/file_worker.py

Lines changed: 37 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -339,22 +339,23 @@ def _process_test(
339339
f"uploading test chunk {chunk.index}/{len(requests)} "
340340
f"({chunk.size_bytes} bytes)",
341341
)
342-
result = transport.post_json(
342+
http_result = transport.post_json(
343343
runtime.endpoints.test_url,
344344
headers,
345345
request.body,
346346
content_encoding=request.content_encoding,
347347
)
348-
requests_attempted += result.attempts
349-
retries += result.retries
348+
requests_attempted += http_result.attempts
349+
retries += http_result.retries
350350
_debug(
351351
runtime,
352352
task,
353-
f"test chunk {chunk.index} completed after {result.attempts} attempt(s)",
353+
f"test chunk {chunk.index} completed after "
354+
f"{http_result.attempts} attempt(s)",
354355
)
355-
if not result.succeeded:
356+
if not http_result.succeeded:
356357
failure_code, failure_message = _http_failure(
357-
result,
358+
http_result,
358359
payload_type=task.payload_type,
359360
payload_limit_context=(
360361
f"chunk={chunk.index}/{len(requests)} "
@@ -407,16 +408,18 @@ def _enrich_and_validate_test(
407408
warnings.append(sidecar_warning)
408409
repo_key = payload_repo_key(bazel_metadata)
409410
context_selection = runtime.context_plan.select(repo_key)
411+
if sidecar_warning:
412+
sidecar_state = "invalid"
413+
elif bazel_metadata is not None:
414+
sidecar_state = "loaded"
415+
else:
416+
sidecar_state = "absent"
410417
_debug(
411418
runtime,
412419
task,
413420
"Bazel sidecar=%s context_repo=%s context_selected=%s context_warning=%s"
414421
% (
415-
"invalid"
416-
if sidecar_warning
417-
else "loaded"
418-
if bazel_metadata is not None
419-
else "absent",
422+
sidecar_state,
420423
repr(repo_key[:256]) if repo_key is not None else "none",
421424
"yes" if context_selection.values is not None else "no",
422425
context_selection.warning_code or "none",
@@ -621,26 +624,26 @@ def _process_coverage(
621624
)
622625
return FileResult(status=FileStatus.SUCCEEDED, **coverage_result_fields)
623626

624-
result = transport.post_prepared_multipart(
627+
http_result = transport.post_prepared_multipart(
625628
runtime.endpoints.coverage_url,
626629
headers,
627630
prepared,
628631
)
629632
_debug(
630633
runtime,
631634
task,
632-
f"coverage request completed after {result.attempts} attempt(s)",
635+
f"coverage request completed after {http_result.attempts} attempt(s)",
633636
)
634-
if not result.succeeded:
637+
if not http_result.succeeded:
635638
failure_code, failure_message = _http_failure(
636-
result,
639+
http_result,
637640
payload_type=task.payload_type,
638641
)
639642
return FileResult(
640643
status=FileStatus.FAILED,
641-
requests_attempted=result.attempts,
644+
requests_attempted=http_result.attempts,
642645
requests_failed=1,
643-
retries=result.retries,
646+
retries=http_result.retries,
644647
failure_code=failure_code,
645648
failure_message=failure_message,
646649
**coverage_result_fields,
@@ -650,9 +653,9 @@ def _process_coverage(
650653
_debug(runtime, task, f"coverage cleanup completed source_deleted={deleted}")
651654
return FileResult(
652655
status=FileStatus.SUCCEEDED,
653-
requests_attempted=result.attempts,
656+
requests_attempted=http_result.attempts,
654657
requests_succeeded=1,
655-
retries=result.retries,
658+
retries=http_result.retries,
656659
source_deleted=deleted,
657660
warning_codes=(cleanup_warning,) if cleanup_warning else (),
658661
**coverage_result_fields,
@@ -702,7 +705,9 @@ def _process_telemetry(
702705
failure_message=metadata_failure,
703706
)
704707

705-
requests = [_TelemetryRequest(source_body, _telemetry_headers(runtime, payload))]
708+
requests: list[_TelemetryRequest] = [
709+
_TelemetryRequest(source_body, _telemetry_headers(runtime, payload))
710+
]
706711
if directive.create_synthetic:
707712
synthetic = _build_synthetic_telemetry(payload, directive, warnings)
708713
if synthetic is not None:
@@ -758,21 +763,22 @@ def _process_telemetry(
758763
requests_succeeded = 0
759764
retries = 0
760765
for request in requests:
761-
result = transport.post_json(
766+
http_result = transport.post_json(
762767
runtime.endpoints.telemetry_url,
763768
request.headers,
764769
request.body,
765770
)
766-
requests_attempted += result.attempts
767-
retries += result.retries
771+
requests_attempted += http_result.attempts
772+
retries += http_result.retries
768773
_debug(
769774
runtime,
770775
task,
771-
f"telemetry request completed after {result.attempts} attempt(s)",
776+
"telemetry request completed after "
777+
f"{http_result.attempts} attempt(s)",
772778
)
773-
if not result.succeeded:
779+
if not http_result.succeeded:
774780
failure_code, failure_message = _http_failure(
775-
result,
781+
http_result,
776782
payload_type=task.payload_type,
777783
)
778784
return FileResult(
@@ -1047,12 +1053,12 @@ def _cleanup_source(path: Path, keep_payloads: bool) -> tuple[bool, str | None]:
10471053

10481054

10491055
def _http_failure(
1050-
result: HttpResult,
1056+
http_result: HttpResult,
10511057
*,
10521058
payload_type: PayloadType,
10531059
payload_limit_context: str | None = None,
10541060
) -> tuple[str, str]:
1055-
if result.status_code == 413:
1061+
if http_result.status_code == 413:
10561062
if payload_type is PayloadType.TEST and payload_limit_context:
10571063
return (
10581064
"payload_limit_contract_mismatch",
@@ -1063,9 +1069,9 @@ def _http_failure(
10631069
f"HTTP 413 for unsplit {payload_type.value} payload; "
10641070
f"{payload_type.value} splitting is not supported",
10651071
)
1066-
if result.status_code is not None:
1067-
return "upload_http_error", f"HTTP {result.status_code}"
1068-
return "upload_transport_error", result.transport_error or "transport error"
1072+
if http_result.status_code is not None:
1073+
return "upload_http_error", f"HTTP {http_result.status_code}"
1074+
return "upload_transport_error", http_result.transport_error or "transport error"
10691075

10701076

10711077
def _event_count(payload: Mapping[str, Any]) -> int:

tools/core/uploader_py/freshness.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@
2828
from .config import UploaderConfig
2929
from .discovery import DiscoveryResult, ScanRoot
3030
from .json_utils import strict_json_loads
31-
from .models import FileResult, FileStatus, PayloadType
31+
from .models import FileResult, FileStatus, FileTask, PayloadType
3232

3333

3434
class FreshnessError(RuntimeError):
@@ -251,7 +251,7 @@ def filter_discovery_for_freshness(
251251
selected_outputs = []
252252
task_label_by_output: dict[str, str] = {}
253253
skipped: list[str] = []
254-
tasks_by_output: dict[str, list[object]] = {}
254+
tasks_by_output: dict[str, list[FileTask]] = {}
255255
for task in discovery.tasks:
256256
tasks_by_output.setdefault(task.output_key or "", []).append(task)
257257

tools/core/uploader_py/reporting.py

Lines changed: 15 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -182,20 +182,20 @@ def statistics(self) -> dict[str, Any]:
182182
"source_files_split": sum(
183183
int(result.chunks_created > 1) for result in self.results
184184
),
185-
"chunks_created": _sum(self.results, "chunks_created"),
186-
"chunks_uploaded": _sum(self.results, "chunks_uploaded"),
187-
"chunks_failed": _sum(self.results, "chunks_failed"),
185+
"chunks_created": sum(result.chunks_created for result in self.results),
186+
"chunks_uploaded": sum(result.chunks_uploaded for result in self.results),
187+
"chunks_failed": sum(result.chunks_failed for result in self.results),
188188
"oversized_single_events": sum(
189189
int(result.failure_code == "single_event_exceeds_payload_limit")
190190
for result in self.results
191191
),
192192
},
193193
"requests": {
194-
"planned": _sum(self.results, "requests_planned"),
195-
"attempted": _sum(self.results, "requests_attempted"),
196-
"succeeded": _sum(self.results, "requests_succeeded"),
197-
"failed": _sum(self.results, "requests_failed"),
198-
"retries": _sum(self.results, "retries"),
194+
"planned": sum(result.requests_planned for result in self.results),
195+
"attempted": sum(result.requests_attempted for result in self.results),
196+
"succeeded": sum(result.requests_succeeded for result in self.results),
197+
"failed": sum(result.requests_failed for result in self.results),
198+
"retries": sum(result.retries for result in self.results),
199199
},
200200
"warnings": dict(sorted(warning_codes.items())),
201201
"failures": dict(sorted(failure_codes.items())),
@@ -211,8 +211,8 @@ def human_lines(self) -> tuple[str, ...]:
211211
concurrency = stats["concurrency"]
212212

213213
def type_value(name: str) -> str:
214-
values = payload_types[name]
215-
return f"{values['succeeded']}/{values['failed']}/{values['skipped']}"
214+
counts = payload_types[name]
215+
return f"{counts['succeeded']}/{counts['failed']}/{counts['skipped']}"
216216

217217
return (
218218
(
@@ -438,7 +438,7 @@ def _legacy_payload_type_counts(
438438
if result.status is FileStatus.SUCCEEDED
439439
)
440440
else:
441-
telemetry_processed = _sum(telemetry, "requests_succeeded")
441+
telemetry_processed = sum(result.requests_succeeded for result in telemetry)
442442
telemetry_failed = sum(
443443
(
444444
result.requests_failed
@@ -501,16 +501,14 @@ def _result_reason(
501501
def _payload_type_counts(
502502
results: tuple[FileResult, ...], payload_type: PayloadType
503503
) -> dict[str, int]:
504-
selected = tuple(result for result in results if result.payload_type is payload_type)
504+
type_results = tuple(
505+
result for result in results if result.payload_type is payload_type
506+
)
505507
return {
506-
status.value: _status_count(selected, status)
508+
status.value: _status_count(type_results, status)
507509
for status in FileStatus
508510
}
509511

510512

511513
def _status_count(results: Iterable[FileResult], status: FileStatus) -> int:
512514
return sum(int(result.status is status) for result in results)
513-
514-
515-
def _sum(results: Iterable[FileResult], field_name: str) -> int:
516-
return sum(int(getattr(result, field_name)) for result in results)

tools/core/uploader_py/resources.py

Lines changed: 25 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -58,44 +58,46 @@ def load_resources(
5858
"""Resolve contexts, telemetry facts, and schema once with warning fallback."""
5959
warnings: list[str] = []
6060
override_path = _existing_file(inputs.context_override)
61-
override_context = _load_json_object(override_path) if override_path else None
62-
if inputs.context_override is not None and override_context is None:
61+
override_values = _load_json_object(override_path) if override_path else None
62+
if inputs.context_override is not None and override_values is None:
6363
warnings.append("context_override_invalid")
6464

6565
context_records: list[ContextRecord] = []
66-
primary_path: Path | None = None
67-
runtime_override_enabled = override_context is not None
68-
if override_context is not None:
69-
context_records.append(ContextRecord.create("__runtime_override__", override_context))
70-
primary_path = override_path
66+
primary_context_path: Path | None = None
67+
runtime_override_enabled = override_values is not None
68+
if override_values is not None:
69+
context_records.append(
70+
ContextRecord.create("__runtime_override__", override_values)
71+
)
72+
primary_context_path = override_path
7173
else:
72-
manifest = _resolve_optional(resolver, inputs.context_manifest_paths)
73-
if inputs.context_manifest_paths and manifest is None:
74+
context_manifest = _resolve_optional(resolver, inputs.context_manifest_paths)
75+
if inputs.context_manifest_paths and context_manifest is None:
7476
warnings.append("context_manifest_unresolved")
75-
if manifest is not None:
77+
if context_manifest is not None:
7678
seen_repo_keys: set[str] = set()
7779
for repo_key, short_path, artifact_path in _context_manifest_entries(
78-
manifest,
80+
context_manifest,
7981
warnings,
8082
):
8183
normalized_repo_key = repo_key.rsplit("+", 1)[-1]
8284
if normalized_repo_key in seen_repo_keys:
8385
warnings.append("context_manifest_duplicate_repo")
8486
continue
85-
resolved = _resolve_optional(
87+
context_path = _resolve_optional(
8688
resolver,
8789
(artifact_path, short_path),
8890
)
89-
context = _load_json_object(resolved) if resolved else None
90-
if context is None:
91+
context_values = _load_json_object(context_path) if context_path else None
92+
if context_values is None:
9193
warnings.append("context_entry_invalid")
9294
continue
9395
seen_repo_keys.add(normalized_repo_key)
9496
context_records.append(
95-
ContextRecord.create(normalized_repo_key, context)
97+
ContextRecord.create(normalized_repo_key, context_values)
9698
)
97-
if primary_path is None:
98-
primary_path = resolved
99+
if primary_context_path is None:
100+
primary_context_path = context_path
99101

100102
primary_context_record = context_records[0] if context_records else None
101103
primary_context = (
@@ -120,13 +122,13 @@ def load_resources(
120122
warnings,
121123
"telemetry_facts_manifest_invalid",
122124
):
123-
resolved = _resolve_optional(resolver, (artifact_path, short_path))
124-
if resolved is None:
125+
facts_path = _resolve_optional(resolver, (artifact_path, short_path))
126+
if facts_path is None:
125127
warnings.append("telemetry_facts_entry_unresolved")
126128
continue
127-
telemetry_facts_paths.append(resolved)
128-
if runtime_override_enabled and primary_path is not None:
129-
sibling = primary_path.parent / "telemetry_facts.json"
129+
telemetry_facts_paths.append(facts_path)
130+
if runtime_override_enabled and primary_context_path is not None:
131+
sibling = primary_context_path.parent / "telemetry_facts.json"
130132
if sibling.is_file():
131133
telemetry_facts_paths.append(sibling.resolve())
132134
telemetry_facts_paths = list(dict.fromkeys(sorted(telemetry_facts_paths)))
@@ -138,7 +140,7 @@ def load_resources(
138140
return LoadedResources(
139141
context_plan=context_plan,
140142
primary_context=primary_context,
141-
primary_context_path=primary_path,
143+
primary_context_path=primary_context_path,
142144
telemetry_facts_paths=tuple(telemetry_facts_paths),
143145
schema=schema,
144146
warning_codes=tuple(dict.fromkeys(warnings)),

0 commit comments

Comments
 (0)