diff --git a/.github/workflows/README.md b/.github/workflows/README.md index fc3a1396..c9605f7a 100644 --- a/.github/workflows/README.md +++ b/.github/workflows/README.md @@ -32,7 +32,7 @@ This directory contains GitHub Actions workflows for automated testing. - Python API tests (18) - Market package platform and job-stream tests (12) - gRPC client/server tests (2) -- Append-job tests (20 C++ + 5 optional gRPC checks) +- Append-job tests (23 C++ + 5 optional gRPC checks) - FCFS/EASY backfill-window focused rerun of the five-check gRPC binary - Synchronized single-coordinator gRPC test - Progressive-loading tests (C++ + CLI) @@ -90,7 +90,7 @@ Total tests referenced by the full suite: | Python API | 18 | CI runner | | Market package | 12 | CI runner; platforms, jobs and trace preparation | | gRPC client/server | 2 | CI runner | -| Append-job | 25: 20 C++ + 5 optional gRPC checks | CI runner | +| Append-job | 28: 23 C++ + 5 optional gRPC checks | CI runner | | FCFS/EASY backfill-window gRPC | 5 repeated checks; 1 targeted | CI runner | | Single-coordinator gRPC | 1 | CI runner; synchronized independent systems | | Progressive loading | 11 C++ + 5 CLI | CI runner | diff --git a/.github/workflows/quick-test.yml b/.github/workflows/quick-test.yml index b9553ac5..55f17b24 100644 --- a/.github/workflows/quick-test.yml +++ b/.github/workflows/quick-test.yml @@ -63,6 +63,10 @@ jobs: $GITHUB_WORKSPACE/install/bin/tests/test_progressive_load $GITHUB_WORKSPACE/install/bin/tests/test_queue_input + - name: Run Streaming Append API Tests + id: append-job + run: $GITHUB_WORKSPACE/install/bin/tests/test_append_job_api + - name: Run Ser20 Memory Serialization Test id: serialization run: $GITHUB_WORKSPACE/install/bin/tests/t_state_ser20 @@ -73,6 +77,7 @@ jobs: echo "========================================" echo "Quick Test Complete" echo "Scheduler correctness (34 fixtures): ${{ steps.scheduler-correctness.outcome }}" + echo "Streaming append API (23 checks): ${{ steps.append-job.outcome }}" echo "Ser20 memory serialization: ${{ steps.serialization.outcome }}" echo "Trace schema / progressive-loading API: ${{ steps.trace-schema.outcome }}" echo "========================================" diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 303df0b9..1d6cee6e 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -287,7 +287,7 @@ jobs: id: append-job run: | echo "========================================" - echo "Running Append-Job Tests (20 C++ + 5 optional gRPC checks)" + echo "Running Append-Job Tests (23 C++ + 5 optional gRPC checks)" echo "========================================" ./tests/run_append_job_tests.sh @@ -352,7 +352,7 @@ jobs: echo " CTest (17 native + trace tools + MPI when available): ${{ steps.core-ctest.outcome }}" echo " Python API (19 tests): ${{ steps.python-api.outcome }}" echo " Market Package (12 tests): ${{ steps.market-package.outcome }}" - echo " Append-Job (20 C++ + 5 optional gRPC checks): ${{ steps.append-job.outcome }}" + echo " Append-Job (23 C++ + 5 optional gRPC checks): ${{ steps.append-job.outcome }}" echo " Backfill-window focused rerun (5 checks already counted): ${{ steps.backfill-window.outcome }}" echo " Single-Coordinator gRPC: ${{ steps.grpc-single-coordinator.outcome }}" echo " Progressive Loading (11 C++ + 5 CLI): ${{ steps.progressive-load.outcome }}" diff --git a/docs/api/PYTHON_API.md b/docs/api/PYTHON_API.md index 314d155e..701e185b 100644 --- a/docs/api/PYTHON_API.md +++ b/docs/api/PYTHON_API.md @@ -98,7 +98,7 @@ Their scheduling semantics are documented in |---|---| | `run()` | Run a configured batch trace to completion. | | `initialize_trace(max_jobs=0)` | Load the configured trace and return the number loaded. | -| `append_job(submit_time, num_nodes, queue, limit_time, actual_run_time=None)` | Append and enqueue one live job; return its ID. A known runtime must be positive and no greater than the limit. | +| `append_job(submit_time, num_nodes, queue, limit_time, actual_run_time=None)` | Append and enqueue one live job; return its ID. The limit must be a positive whole number of seconds; a known runtime must be positive and no greater than the limit. | | `append_jobs(requests)` | Atomically append and enqueue ordered `JobAppendRequest` values; return their IDs. | | `get_job_statuses(job_idxs)` | Return lifecycle and timing snapshots for appended job IDs. | | `advance_to(target_time)` | Process events at or before the target. | diff --git a/docs/api/STREAMING_API.md b/docs/api/STREAMING_API.md index 217e4882..2394d99a 100644 --- a/docs/api/STREAMING_API.md +++ b/docs/api/STREAMING_API.md @@ -70,7 +70,8 @@ job_no_t append_job(sim_time_t submit_time, num_nodes_t num_nodes, - `num_nodes`: Number of nodes the job requests - `queue`: Numeric queue ID (for example, `"1"`) in the default build, or a queue name (for example, `"pbatch"`) with `DR_EVT_LEGACY_QUEUE_INPUT` -- `limit_time`: User-estimated time limit, in seconds +- `limit_time`: Positive whole-number time limit, in seconds. Fractional, + non-finite, zero, and negative values are rejected. - `actual_run_time`: Optional known execution duration. It must be positive and no greater than `limit_time`. When omitted, streaming jobs retain the existing behavior of running for `limit_time`. @@ -93,6 +94,7 @@ jobs are already known together (e.g. several arrivals collected in one polling interval), not just a loop over `append_job()`. All-or-nothing: requests must already be sorted by `submit_time` (non-decreasing), and either the whole batch is appended or, on any failure (unsorted input, +an invalid or fractional `limit_time`, `--job_store_overflow=abort` with no room even after reclaiming, or `--check_memory_pressure` refusing the batch under real memory pressure - see [Command-Line Options](../user-guide/command-line.md)), diff --git a/docs/user-guide/trace-formats.md b/docs/user-guide/trace-formats.md index 3a2ea4c9..b2b1750c 100644 --- a/docs/user-guide/trace-formats.md +++ b/docs/user-guide/trace-formats.md @@ -92,7 +92,7 @@ determines simulation vs replay mode (see below). | `job_submit_time` | When the job arrives/submits. Accepted alias: `submit_time` | Both modes | | `num_nodes` | Number of nodes requested | Both modes | | `q_id` | Optional one-based queue ID. If absent, the job uses `1` (`Queue1`). | Both modes | -| `time_limit` | User-provided time limit (seconds). Accepted column-name aliases: `time_limit`, `timelimit`, `walltime` | Both modes | +| `time_limit` | Positive whole-number time limit (seconds); fractional values are rejected. Accepted column-name aliases: `time_limit`, `timelimit`, `walltime` | Both modes | | `begin_time` | Historical start time of this individual job; distinct from the global `--sim_start_time` boundary | Replay mode only; must appear together with `end_time` | | `end_time` | Historical end time from trace | Replay mode only; must appear together with `begin_time` | | `avgpcon` | Average power usage associated with the job | Required only with `--trace_type pcon`; ignored in standard mode | diff --git a/experimental/multi-cluster/README.md b/experimental/multi-cluster/README.md index 719fe586..ccbea58f 100644 --- a/experimental/multi-cluster/README.md +++ b/experimental/multi-cluster/README.md @@ -34,7 +34,7 @@ submit_time,num_nodes,time_limit,duration |---|---| | `submit_time` | Arrival time in simulation seconds. Rows must be in nondecreasing order. The legacy alias `job_submit_time` is accepted. | | `num_nodes` | Positive node request. | -| `time_limit` | Positive requested wall-time limit. | +| `time_limit` | Positive whole-number wall-time limit in seconds; fractional values are rejected. | | `duration` | Positive reference duration, no greater than `time_limit`. The legacy alias `actual_run_time` is accepted. | `job_id` is optional and defaults to the zero-based row number. The queue is @@ -173,6 +173,9 @@ python3 experimental/multi-cluster/plot_relative_performance.py \ The default logarithmic axes make both sub-unit and large speedups visible. Use `--linear` for linear axes or `--systems tuolumne tuolumne-cpu dane` to select panels. Unsuffixed GPU-machine names resolve to their `-gpu` columns. +Cells lacking either an actual or predicted value are reported and omitted +because they cannot form an actual-versus-predicted point. When `--output` +does not already name a PDF, the plotter also writes a PDF with the same stem. ## Sampling @@ -218,7 +221,7 @@ The experiment assumes that the submitting user can correct an underestimated wall-time request using the known ground-truth runtime before the successful submission. Starting from `predicted_time_limit`, it doubles the request until it is at least `actual_duration`, without simulating or charging resources for -the failed attempts. `--max-time-limit` caps the corrected request; the +the failed attempts. `--max-time-limit` caps the adapted request; the prediction study defaults to Lassen's 43,200-second maximum. If `actual_duration` exceeds that maximum, the experiment writes a `dropped:` record to standard error and does not submit the job. Dropped jobs are excluded @@ -240,10 +243,11 @@ capacity-compatible system can finish within the maximum, the job is dropped. it chooses the highest predicted relative performance over all feasible systems and lets that system queue the job. Ties use systems-table order. -`--wall-time-policy corrected-prediction` uses the doubling behavior described -above. `--wall-time-policy actual-duration` instead submits a time limit exactly -equal to the selected system's ground-truth runtime. This second policy never -changes or truncates the runtime; it changes only the requested time limit. +`--wall-time-policy adapted-limit` uses the doubling behavior described +above. `--wall-time-policy actual-duration` instead submits the smallest +whole-second time limit that covers the selected system's ground-truth runtime. +This second policy never changes or truncates the runtime; it changes only the +requested time limit. The wait estimate uses the submitted time limit. Only the selected worker receives the job, with that limit and the actual duration. Thus prediction @@ -267,7 +271,7 @@ mpirun -np 6 build/mpi_performance_dispatch \ --seed 7 \ --max-time-limit 43200 \ --dispatch-policy IPDPS24 \ - --wall-time-policy corrected-prediction \ + --wall-time-policy adapted-limit \ --output dispatch-decisions.csv ``` @@ -277,7 +281,7 @@ allocation. MPI and Ser20 are required; gRPC and Python are not. ### Prediction-study matrix The prediction-study runner executes both dispatch policies (`turnaround` and -`IPDPS24`) with both wall-time policies (`corrected-prediction` and +`IPDPS24`) with both wall-time policies (`adapted-limit` and `actual-duration`) for the ideal, application-average, and RAJAPerf prediction tables. Each of these 12 configurations runs on all ten synthetic traces, for 120 runs by default: @@ -301,6 +305,13 @@ Output filenames and completion markers include the dispatch and wall-time policy names. This prevents results from different configurations from being mistaken for one another. +After all configurations are available, `summary.csv` contains one row for +each of the 120 runs, while `summary_aggregate.csv` and `summary.md` contain +the 12 ten-trace mean/standard-deviation records. `summary.png` plots those +aggregate metrics in four panels, and `summary.pdf` contains the same figure +in vector form. `metrics_per_run.csv` is retained as a compatibility copy of +`summary.csv`. + ## Output Rank 0 writes one CSV row per successfully dispatched job; dropped jobs are @@ -322,8 +333,8 @@ recorded in the log instead: | `actual_duration` | Ground-truth duration submitted to DR_EVT. | | `predicted_time_limit` | Wall-time limit used to estimate the dispatch candidate. | | `actual_time_limit` | Ground-truth-scaled wall-time limit retained for comparison. | -| `submitted_time_limit` | Limit submitted to DR_EVT: the corrected predicted limit or the exact ground-truth runtime, according to `--wall-time-policy`. | -| `time_limit_doublings` | Number of pre-submission doublings needed by `corrected-prediction`; zero for `actual-duration`. | +| `submitted_time_limit` | Whole-second limit submitted to DR_EVT: the upward-rounded adapted predicted limit or ground-truth runtime, according to `--wall-time-policy`. | +| `time_limit_doublings` | Number of pre-submission doublings needed by the adapted predicted-limit policy; zero for `actual-duration`. | | `predicted_turnaround` | `estimated_wait + estimated_duration`. | | `job_idx` | Worker-local DR_EVT job identifier. | @@ -363,7 +374,7 @@ python3 python/grpc_mpi_launcher.py --mpi-ranks 6 \ --seed 7 \ --max-time-limit 43200 \ --dispatch-policy IPDPS24 \ - --wall-time-policy corrected-prediction \ + --wall-time-policy adapted-limit \ --output grpc-dispatch-decisions.csv ``` diff --git a/experimental/multi-cluster/grpc_performance_dispatch.py b/experimental/multi-cluster/grpc_performance_dispatch.py index 178264c4..b72ccd00 100644 --- a/experimental/multi-cluster/grpc_performance_dispatch.py +++ b/experimental/multi-cluster/grpc_performance_dispatch.py @@ -333,12 +333,12 @@ def choose_system( horizons, max_time_limit=math.inf, dispatch_policy="turnaround", - wall_time_policy="corrected-prediction", + wall_time_policy="adapted-limit", ): """Choose a feasible system using turnaround or paper Algorithm 2.""" if dispatch_policy not in {"turnaround", "IPDPS24"}: raise ValueError(f"unknown dispatch policy: {dispatch_policy}") - if wall_time_policy not in {"corrected-prediction", "actual-duration"}: + if wall_time_policy not in {"adapted-limit", "actual-duration"}: raise ValueError(f"unknown wall-time policy: {wall_time_policy}") candidates = [] for index, (system, window, horizon) in enumerate( @@ -366,6 +366,9 @@ def choose_system( submitted_limit, doublings = adjusted_time_limit( predicted_limit, actual_duration, max_time_limit ) + submitted_limit = math.ceil(submitted_limit) + if submitted_limit > max_time_limit: + continue wait = estimate_wait(window, job["num_nodes"], submitted_limit, horizon) if math.isfinite(wait) or dispatch_policy == "IPDPS24": candidates.append( @@ -716,11 +719,11 @@ def main(): ) parser.add_argument( "--wall-time-policy", - choices=("corrected-prediction", "actual-duration"), - default="corrected-prediction", + choices=("adapted-limit", "actual-duration"), + default="adapted-limit", help=( - "submit a corrected predicted limit or the ground-truth runtime " - "(default: corrected-prediction)" + "submit an adapted predicted limit or the ground-truth runtime " + "(default: adapted-limit)" ), ) parser.add_argument("--session-name", default="performance-dispatch") diff --git a/experimental/multi-cluster/mpi_performance_dispatch.cpp b/experimental/multi-cluster/mpi_performance_dispatch.cpp index 09b97197..efeb187b 100644 --- a/experimental/multi-cluster/mpi_performance_dispatch.cpp +++ b/experimental/multi-cluster/mpi_performance_dispatch.cpp @@ -58,7 +58,7 @@ constexpr int kPayloadTag = 101; enum class Operation : std::uint8_t { Snapshot, Append, Finish }; enum class DispatchPolicy : std::uint8_t { Turnaround, Ipdps24 }; enum class WallTimePolicy : std::uint8_t { - CorrectedPrediction, + AdaptedLimit, ActualDuration }; @@ -171,7 +171,7 @@ struct Options { double prediction_utilization = 1.0; double max_time_limit = std::numeric_limits::infinity(); DispatchPolicy dispatch_policy = DispatchPolicy::Turnaround; - WallTimePolicy wall_time_policy = WallTimePolicy::CorrectedPrediction; + WallTimePolicy wall_time_policy = WallTimePolicy::AdaptedLimit; }; [[noreturn]] void usage(const char *program, const std::string &error = {}) { @@ -189,8 +189,8 @@ struct Options { "(default: unlimited)\n" << " --dispatch-policy POLICY turnaround or IPDPS24 " "(default: turnaround)\n" - << " --wall-time-policy POLICY corrected-prediction or " - "actual-duration (default: corrected-prediction)\n" + << " --wall-time-policy POLICY adapted-limit or " + "actual-duration (default: adapted-limit)\n" << " --output PATH decision CSV (default: stdout)\n"; throw std::invalid_argument(error.empty() ? "help requested" : error); } @@ -233,12 +233,12 @@ Options parse_options(int argc, char **argv) { usage(argv[0], "--dispatch-policy must be turnaround or IPDPS24"); } else if (arg == "--wall-time-policy") { const auto value = option_value(i, argc, argv, arg); - if (value == "corrected-prediction") - options.wall_time_policy = WallTimePolicy::CorrectedPrediction; + if (value == "adapted-limit") + options.wall_time_policy = WallTimePolicy::AdaptedLimit; else if (value == "actual-duration") options.wall_time_policy = WallTimePolicy::ActualDuration; else - usage(argv[0], "--wall-time-policy must be corrected-prediction or " + usage(argv[0], "--wall-time-policy must be adapted-limit or " "actual-duration"); } else if (arg == "--output") @@ -826,11 +826,14 @@ choose_system(const Job &job, const Workload &workload, continue; const double predicted_limit = job.limit_time / performance->predicted; const double actual_limit = job.limit_time / performance->ground_truth; - const auto [submitted_limit, doublings] = + auto [submitted_limit, doublings] = wall_time_policy == WallTimePolicy::ActualDuration ? std::pair{actual_duration, std::uint32_t{0}} : adjust_time_limit(predicted_limit, actual_duration, max_time_limit); + submitted_limit = std::ceil(submitted_limit); + if (submitted_limit > max_time_limit) + continue; const double wait = estimate_wait(snapshots[i].window, job.num_nodes, submitted_limit, snapshots[i].prediction_horizon); diff --git a/experimental/multi-cluster/plot_relative_performance.py b/experimental/multi-cluster/plot_relative_performance.py index 1b69fc6b..6d4a0a72 100644 --- a/experimental/multi-cluster/plot_relative_performance.py +++ b/experimental/multi-cluster/plot_relative_performance.py @@ -5,6 +5,7 @@ import csv import math import pathlib +import sys IDENTITY_COLUMNS = {"App", "Args", "Ranks", "app", "args", "ranks"} @@ -45,6 +46,7 @@ def load_points(ground_truth_path, prediction_path, requested_systems=None): if len(set(systems)) != len(systems): raise ValueError("the requested systems resolve to duplicates") points = {system: [] for system in systems} + skipped = {system: 0 for system in systems} for row_number, (row, predicted_row) in enumerate( zip(actual_reader, predicted_reader, strict=True), start=2 ): @@ -60,9 +62,8 @@ def load_points(ground_truth_path, prediction_path, requested_systems=None): if not actual_text and not predicted_text: continue if not actual_text or not predicted_text: - raise ValueError( - f"row {row_number}: incomplete {system} pair" - ) + skipped[system] += 1 + continue actual = float(actual_text) predicted = float(predicted_text) if ( @@ -73,6 +74,15 @@ def load_points(ground_truth_path, prediction_path, requested_systems=None): ): raise ValueError(f"row {row_number}: invalid {system} pair") points[system].append((actual, predicted, app)) + skipped = {system: count for system, count in skipped.items() if count} + if skipped: + detail = ", ".join( + f"{system}={count}" for system, count in skipped.items() + ) + print( + f"warning: skipped unpaired actual/predicted cells: {detail}", + file=sys.stderr, + ) return points @@ -176,6 +186,8 @@ def plot_points(points, output, logarithmic=True): figure.tight_layout(rect=(0, 0.08, 1, 0.96)) output.parent.mkdir(parents=True, exist_ok=True) figure.savefig(output, dpi=180, bbox_inches="tight") + if output.suffix.lower() != ".pdf": + figure.savefig(output.with_suffix(".pdf"), bbox_inches="tight") plt.close(figure) diff --git a/experimental/multi-cluster/run_prediction_study.py b/experimental/multi-cluster/run_prediction_study.py index 703ed4b8..125438f0 100644 --- a/experimental/multi-cluster/run_prediction_study.py +++ b/experimental/multi-cluster/run_prediction_study.py @@ -3,7 +3,7 @@ The study uses ten Lassen synthetic job streams and compares ideal, application-average, and RAJAPerf predictions under the turnaround-aware and -IPDPS24 dispatch policies. Each combination is run with corrected predicted +IPDPS24 dispatch policies. Each combination is run with adapted predicted wall times and with wall time set to actual duration. A model-based case can be added by passing ``--model-prediction``. """ @@ -25,7 +25,8 @@ "average_speedup", ) DISPATCH_POLICIES = ("turnaround", "IPDPS24") -WALL_TIME_POLICIES = ("corrected-prediction", "actual-duration") +WALL_TIME_POLICIES = ("adapted-limit", "actual-duration") +JOBS_PER_TRACE = 100_000 OVERALL_RE = re.compile( r"^overall: jobs=(?P\d+) dropped_jobs=(?P\d+) " + r" ".join(rf"{name}=(?P<{name}>[-+0-9.eE]+)" for name in METRICS) @@ -81,8 +82,10 @@ def validate_inputs(root, executable, cases): for trace in traces: with trace.open(encoding="utf-8") as stream: rows = sum(1 for _ in stream) - 1 - if rows != 100_000: - raise ValueError(f"{trace} contains {rows} jobs, expected 100000") + if rows != JOBS_PER_TRACE: + raise ValueError( + f"{trace} contains {rows} jobs, expected {JOBS_PER_TRACE}" + ) return traces @@ -118,7 +121,7 @@ def run_case( rows = sum(1 for _ in stream) - 1 if ( rows == record["jobs"] - and rows + record["dropped_jobs"] == 100_000 + and rows + record["dropped_jobs"] == JOBS_PER_TRACE ): print(f"skip validated {stem}", flush=True) return record @@ -162,7 +165,10 @@ def run_case( record = parse_overall(combined) with dispatch.open(encoding="utf-8") as stream: rows = sum(1 for _ in stream) - 1 - if rows != record["jobs"] or rows + record["dropped_jobs"] != 100_000: + if ( + rows != record["jobs"] + or rows + record["dropped_jobs"] != JOBS_PER_TRACE + ): raise RuntimeError( f"{stem} is incomplete: dispatch rows={rows}, " f"reported jobs={record['jobs']}, dropped={record['dropped_jobs']}" @@ -173,23 +179,24 @@ def run_case( def write_results(output_dir, records): """Write per-run data and ten-run aggregate tables.""" - per_run = output_dir / "metrics_per_run.csv" - with per_run.open("w", newline="", encoding="utf-8") as stream: - writer = csv.DictWriter( - stream, - fieldnames=( - "dispatch_policy", - "wall_time_policy", - "case", - "run", - "jobs", - "dropped_jobs", - *METRICS, - ), - lineterminator="\n", - ) - writer.writeheader() - writer.writerows(records) + per_run_fields = ( + "dispatch_policy", + "wall_time_policy", + "case", + "run", + "jobs", + "dropped_jobs", + *METRICS, + ) + for filename in ("summary.csv", "metrics_per_run.csv"): + with (output_dir / filename).open( + "w", newline="", encoding="utf-8" + ) as stream: + writer = csv.DictWriter( + stream, fieldnames=per_run_fields, lineterminator="\n" + ) + writer.writeheader() + writer.writerows(records) grouped = {} for record in records: @@ -222,7 +229,7 @@ def write_results(output_dir, records): row[f"{metric}_stddev"] = statistics.pstdev(values) summary_rows.append(row) - summary_csv = output_dir / "summary.csv" + summary_csv = output_dir / "summary_aggregate.csv" fields = [ "dispatch_policy", "wall_time_policy", @@ -250,7 +257,7 @@ def write_results(output_dir, records): "", "Values are the mean ± population standard deviation over 10 job traces.", "", - "| Dispatch | Wall time | Prediction | Dropped jobs | Turnaround time | Bounded slowdown | Run time | Speedup |", + "| Dispatch | Wall time | Prediction | Dropped jobs | Turnaround time (sec) | Bounded slowdown | Run time (sec) | Speedup |", "|---|---|---|---:|---:|---:|---:|---:|", ] for row in summary_rows: @@ -287,30 +294,73 @@ def plot_results(output_dir, summary_rows): "matplotlib is required for the plot; tables were still generated" ) from error - labels = [ - f"{row['dispatch_policy']}\n{row['wall_time_policy']}\n" - f"{row['case'].replace('_', ' ')}" + cases = list(dict.fromkeys(row["case"] for row in summary_rows)) + configurations = list( + dict.fromkeys( + (row["dispatch_policy"], row["wall_time_policy"]) + for row in summary_rows + ) + ) + indexed = { + (row["dispatch_policy"], row["wall_time_policy"], row["case"]): row for row in summary_rows - ] + } + case_labels = { + "ideal": "Ideal", + "model": "Model-based", + "app_avg": "Application average", + "rajaperf": "RAJAPerf", + } + configuration_labels = { + ("turnaround", "adapted-limit"): "Turnaround", + ("turnaround", "actual-duration"): "Turnaround / limit = duration", + ("IPDPS24", "adapted-limit"): "IPDPS24", + ("IPDPS24", "actual-duration"): "IPDPS24 / limit = duration", + } titles = { - "average_turnaround_time": "Average turnaround time", + "average_turnaround_time": "Average turnaround time (sec)", "average_bounded_slowdown": "Average bounded slowdown", - "average_run_time": "Average run time", - "average_speedup": "Average selected speedup", + "average_run_time": "Average run time (sec)", + "average_speedup": "Average speedup", } - fig, axes = plt.subplots(2, 2, figsize=(18, 10)) + fig, axes = plt.subplots(2, 2, figsize=(15, 10)) colors = ["#4C78A8", "#F58518", "#54A24B", "#E45756"] - bar_colors = [colors[index % len(colors)] for index in range(len(labels))] + x_positions = list(range(len(cases))) + width = min(0.18, 0.8 / max(len(configurations), 1)) for axis, metric in zip(axes.flat, METRICS): - means = [row[f"{metric}_mean"] for row in summary_rows] - errors = [row[f"{metric}_stddev"] for row in summary_rows] - axis.bar(labels, means, yerr=errors, capsize=4, color=bar_colors) + for index, configuration in enumerate(configurations): + offset = (index - (len(configurations) - 1) / 2) * width + positions = [position + offset for position in x_positions] + rows = [indexed[(*configuration, case)] for case in cases] + axis.bar( + positions, + [row[f"{metric}_mean"] for row in rows], + width, + yerr=[row[f"{metric}_stddev"] for row in rows], + capsize=3, + color=colors[index % len(colors)], + label=configuration_labels.get( + configuration, " / ".join(configuration) + ), + ) axis.set_title(titles[metric]) + axis.set_xticks(x_positions) + axis.set_xticklabels( + [case_labels.get(case, case.replace("_", " ")) for case in cases] + ) axis.grid(axis="y", alpha=0.25) - axis.tick_params(axis="x", rotation=15) fig.suptitle("Multi-cluster prediction policies (mean ± SD, 10 traces)") - fig.tight_layout() + handles, labels = axes.flat[0].get_legend_handles_labels() + fig.legend( + handles, + labels, + loc="outside lower center", + ncol=min(2, len(labels)), + frameon=False, + ) + fig.tight_layout(rect=(0, 0.1, 1, 0.96)) fig.savefig(output_dir / "summary.png", dpi=180) + fig.savefig(output_dir / "summary.pdf") plt.close(fig) @@ -321,27 +371,47 @@ def load_completed(output_dir, cases, dispatch_policies, wall_time_policies): for wall_time_policy in wall_time_policies: for case, _ in cases: for run_number in range(1, 11): - stem = ( + qualified_stem = ( f"{dispatch_policy}.{wall_time_policy}.{case}." f"run_{run_number:02d}" ) - marker = output_dir / f"{stem}.complete" - log = output_dir / f"{stem}.log" - if not marker.is_file() or not log.is_file(): - raise FileNotFoundError(f"missing completed run: {stem}") expected_marker = ( f"dispatch_policy={dispatch_policy}\n" f"wall_time_policy={wall_time_policy}\n" ) + marker = output_dir / f"{qualified_stem}.complete" + log = output_dir / f"{qualified_stem}.log" + dispatch = output_dir / f"{qualified_stem}.dispatch.csv" + if not ( + marker.is_file() and log.is_file() and dispatch.is_file() + ): + raise FileNotFoundError( + f"missing completed run: {qualified_stem}" + ) + if marker.read_text(encoding="utf-8") != expected_marker: - raise ValueError(f"configuration mismatch for {stem}") + raise ValueError( + f"configuration mismatch for {qualified_stem}" + ) + record = parse_overall(log.read_text(encoding="utf-8")) + with dispatch.open(encoding="utf-8") as stream: + rows = sum(1 for _ in stream) - 1 + if ( + rows != record["jobs"] + or rows + record["dropped_jobs"] != JOBS_PER_TRACE + ): + raise ValueError( + f"incomplete run {qualified_stem}: dispatch rows={rows}, " + f"reported jobs={record['jobs']}, " + f"dropped={record['dropped_jobs']}" + ) records.append( { "dispatch_policy": dispatch_policy, "wall_time_policy": wall_time_policy, "case": case, "run": run_number, - **parse_overall(log.read_text()), + **record, } ) return records diff --git a/experimental/multi-cluster/test_performance_dispatch.py b/experimental/multi-cluster/test_performance_dispatch.py index ce5ce909..d2996312 100644 --- a/experimental/multi-cluster/test_performance_dispatch.py +++ b/experimental/multi-cluster/test_performance_dispatch.py @@ -2,6 +2,7 @@ """Unit tests for the Python/gRPC multi-cluster dispatcher policy.""" import io +import math import pathlib import random import tempfile @@ -287,6 +288,22 @@ def test_plot_system_alias_and_point_loading(self): points = performance_plot.load_points(actual, predicted, ["tuolumne"]) self.assertEqual(points, {"tuolumne-gpu": [(1.25, 1.5, "solver")]}) + def test_plot_skips_unpaired_prediction_cells(self): + with tempfile.TemporaryDirectory() as directory: + actual = pathlib.Path(directory) / "actual.csv" + predicted = pathlib.Path(directory) / "predicted.csv" + actual.write_text( + "App,Args,Ranks,dane\nsolver,x,8,\n", encoding="utf-8" + ) + predicted.write_text( + "App,Args,Ranks,dane\nsolver,x,8,1.5\n", encoding="utf-8" + ) + with patch("sys.stderr", new=io.StringIO()) as errors: + points = performance_plot.load_points(actual, predicted) + + self.assertEqual(points, {"dane": []}) + self.assertIn("dane=1", errors.getvalue()) + def test_application_average_baseline_groups_by_app_rank_and_mode(self): rows = [ {"App": "a", "Args": "x", "Ranks": "4", "dane": "1"}, @@ -504,13 +521,13 @@ def test_predicted_time_limit_is_doubled_and_capped(self): self.assertEqual(adjusted_time_limit(10, 120, 100), (100, 4)) self.assertEqual(adjusted_time_limit(40, 40, 100), (40, 0)) - def test_actual_duration_wall_time_is_exact(self): + def test_actual_duration_wall_time_is_rounded_up(self): system = self.systems[0] choice = choose_system( { "job_id": "actual-duration-limit", "num_nodes": 8, - "duration": 80, + "duration": 80.25, "limit_time": 10, }, self.workloads["cpu-solver"][0], @@ -520,7 +537,12 @@ def test_actual_duration_wall_time_is_exact(self): max_time_limit=1000, wall_time_policy="actual-duration", ) - self.assertEqual(choice["submitted_time_limit"], choice["actual_duration"]) + self.assertEqual( + choice["submitted_time_limit"], math.ceil(choice["actual_duration"]) + ) + self.assertGreaterEqual( + choice["submitted_time_limit"], choice["actual_duration"] + ) self.assertEqual(choice["time_limit_doublings"], 0) def test_runtime_limit_filters_systems_before_wait_comparison(self): @@ -576,7 +598,7 @@ def test_controller_uses_same_inputs_and_submits_both_scaled_times(self): prediction_utilization=1.0, max_time_limit=1000.0, dispatch_policy="turnaround", - wall_time_policy="corrected-prediction", + wall_time_policy="adapted-limit", session_name="test", ) FakeSession.instances = [] @@ -724,9 +746,14 @@ def test_prediction_study_aggregates_full_policy_matrix(self): output = pathlib.Path(directory) summary = prediction_study.write_results(output, records) per_run = (output / "metrics_per_run.csv").read_bytes() + complete_summary = (output / "summary.csv").read_bytes() + aggregate_summary = (output / "summary_aggregate.csv").read_text() self.assertEqual(len(summary), 12) self.assertNotIn(b"\r\n", per_run) + self.assertEqual(complete_summary, per_run) + self.assertEqual(len(complete_summary.splitlines()), 121) + self.assertEqual(len(aggregate_summary.splitlines()), 13) self.assertEqual( { (row["dispatch_policy"], row["wall_time_policy"]) @@ -739,6 +766,5 @@ def test_prediction_study_aggregates_full_policy_matrix(self): }, ) - if __name__ == "__main__": unittest.main() diff --git a/src/proto/dr_evt_service.proto b/src/proto/dr_evt_service.proto index b58bce55..00b78c68 100644 --- a/src/proto/dr_evt_service.proto +++ b/src/proto/dr_evt_service.proto @@ -121,7 +121,7 @@ message AppendJobRequest { // Numeric queue ID by default (e.g. "1" for Queue1), or a legacy queue // name (e.g. "pbatch") when DR_EVT_LEGACY_QUEUE_INPUT is enabled. string queue = 3; - double limit_time = 4; // user-estimated time limit, in seconds + double limit_time = 4; // positive whole-second time limit optional double actual_run_time = 5; // known execution time, if available } @@ -132,7 +132,7 @@ message JobAppendData { uint32 num_nodes = 2; // Uses the same compile-time queue representation as AppendJobRequest. string queue = 3; - double limit_time = 4; + double limit_time = 4; // positive whole-second time limit optional double actual_run_time = 5; } diff --git a/src/sim/sim.cpp b/src/sim/sim.cpp index 5af06b91..486afb04 100644 --- a/src/sim/sim.cpp +++ b/src/sim/sim.cpp @@ -24,10 +24,12 @@ #endif #include #include +#include #include #include #include #include +#include #include #include #include @@ -39,6 +41,22 @@ namespace { constexpr tdiff_t bounded_slowdown_threshold = 10.0; +/** Validate and convert a streaming wall-time request. + * @param[in] limit_time Requested wall time in seconds. + * @return The equivalent integral scheduler limit. + * @throws std::invalid_argument if limit_time is non-positive, non-finite, + * fractional, or outside the range of timeout_t. + */ +timeout_t checked_streaming_limit(tdiff_t limit_time) { + if (!std::isfinite(limit_time) || limit_time <= 0.0 || + std::trunc(limit_time) != limit_time || + limit_time > static_cast(std::numeric_limits::max())) { + throw std::invalid_argument( + "limit_time must be a positive whole number of seconds"); + } + return static_cast(limit_time); +} + #if defined(DR_EVT_HAS_SER20) constexpr std::array checkpoint_magic{'D', 'R', 'E', 'V', @@ -1755,9 +1773,10 @@ job_no_t BasicSimulation::append_job(sim_time_t submit_time, std::to_string(submit_time) + " but current_time=" + std::to_string(m_current_time)); } + const timeout_t stored_limit = checked_streaming_limit(limit_time); if (actual_run_time && (!std::isfinite(*actual_run_time) || *actual_run_time <= 0.0 || - *actual_run_time > limit_time)) { + *actual_run_time > static_cast(stored_limit))) { throw std::invalid_argument( "actual_run_time must be positive and no greater than limit_time"); } @@ -1784,10 +1803,9 @@ job_no_t BasicSimulation::append_job(sim_time_t submit_time, #endif job_no_t job_idx = m_trace.append_job(m_current_time, submit_epoch, num_nodes, - q, static_cast(limit_time), - actual_run_time); + q, stored_limit, actual_run_time); submit_job(job_idx, submit_time); - record_appended_job(job_idx, submit_time, num_nodes, limit_time); + record_appended_job(job_idx, submit_time, num_nodes, stored_limit); return job_idx; } @@ -1816,10 +1834,12 @@ std::vector BasicSimulation::append_jobs( " has submit_time=" + std::to_string(requests[i].submit_time) + " but current_time=" + std::to_string(m_current_time)); } + const timeout_t stored_limit = + checked_streaming_limit(requests[i].limit_time); if (requests[i].actual_run_time && (!std::isfinite(*requests[i].actual_run_time) || *requests[i].actual_run_time <= 0.0 || - *requests[i].actual_run_time > requests[i].limit_time)) { + *requests[i].actual_run_time > static_cast(stored_limit))) { throw std::invalid_argument( "request " + std::to_string(i) + " actual_run_time must be positive and no greater than limit_time"); @@ -1848,7 +1868,8 @@ std::vector BasicSimulation::append_jobs( for (size_t i = 0; i < job_idxs.size(); ++i) { submit_job(job_idxs[i], requests[i].submit_time); record_appended_job(job_idxs[i], requests[i].submit_time, - requests[i].num_nodes, requests[i].limit_time); + requests[i].num_nodes, + static_cast(requests[i].limit_time)); } return job_idxs; } diff --git a/src/sim/sim.hpp b/src/sim/sim.hpp index fc6ec60e..3cd0d867 100644 --- a/src/sim/sim.hpp +++ b/src/sim/sim.hpp @@ -265,8 +265,9 @@ template class BasicSimulation { * @param[in] actual_run_time Known execution time, or empty to use * limit_time as the streaming execution duration. * @return New Trace job identifier as job_no_t. - * @throws std::invalid_argument if actual_run_time is non-positive or - * greater than limit_time. + * @throws std::invalid_argument if limit_time is not a positive, + * representable whole number of seconds, or if actual_run_time is + * non-positive, non-finite, or greater than limit_time. * @see submit_job() * @see SchedulerBase::insert_job() */ @@ -290,6 +291,9 @@ template class BasicSimulation { * Trace::append_jobs() for why - this function forwards * requests as-is, so pass them there already sorted). * @return New Trace job identifiers in the same order as requests. + * @throws std::invalid_argument if any limit_time is not a positive, + * representable whole number of seconds, or if an actual_run_time + * is invalid or exceeds its limit_time. * @see append_job() * @see submit_job() */ diff --git a/tests/README.md b/tests/README.md index 6a383992..acfe85cc 100644 --- a/tests/README.md +++ b/tests/README.md @@ -97,13 +97,13 @@ server; no pre-existing Redis service is used. | Job store | 6 | `run_job_store_tests.sh` | Capacity, growth/abort, reclamation, and statistics | | Redis output | 1 integration runner | `run_redis_tests.sh` | Isolated server startup, CSV/hash output, time/resource indexes, namespace replacement, finalized-versus-unfinished visibility, pipelined bulk lookup, and byte-identical job and resource outputs for Redis/file runs of 200 jobs. Redis coverage runs with and without Ser20; the checkpoint portion runs only in the Ser20-enabled configuration. | | Redis gRPC example | 1 integration runner | `run_redis_grpc_client_test.sh` | Batch append, advance, pipelined finalized-job lookup, one server fallback query, and original-order merged reporting | -| Append-job | 26 | `run_append_job_tests.sh` | 21 in-process C++ checks plus 5 optional gRPC checks, including known actual runtimes, capacity-aware instantaneous/aggregate utilization, warm start, and validation | +| Append-job | 28 | `run_append_job_tests.sh` | 23 in-process C++ checks plus 5 optional gRPC checks, including integral streaming-limit validation, known actual runtimes, capacity-aware instantaneous/aggregate utilization, warm start, and validation | | Checkpoint/restart | 7 groups | CTest (`test_checkpoint_restart`) and `run_redis_tests.sh` | Exact scheduler continuation plus byte-identical file, progressive-loading, and Redis job/resource output after archive-and-stitch recovery. | | Progressive loading | 16 | `run_progressive_load_tests.sh` | 11 C++ checks plus 5 CLI checks for multi-file loading, bounded storage, block-queue integration, and memory checks | | Protobuf configuration | 12 | `run_configs_tests.sh` | Configuration/CLI parity, capacity/simulation-start-time validation, and documented examples | | Python API | 19 | `run_python_tests.sh` | Bindings, callbacks, streaming, checkpoint/restart, monitoring, policy APIs, and warm-start execution | | gRPC client/server | 2 | `run_grpc_tests.sh` | Single-pair and optional MPI multi-server behavior | -| Multi-cluster dispatch | 2 | CTest (`test_python_performance_dispatch`, `test_mpi_performance_dispatch`) | Python/gRPC policy and native MPI coverage for app/workload sampling, CPU/GPU compatibility, runtime scaling, turnaround and IPDPS24 placement, corrected and actual-duration wall-time policies, evaluation metrics, machine-size eligibility, and oversized-request truncation | +| Multi-cluster dispatch | 2 | CTest (`test_python_performance_dispatch`, `test_mpi_performance_dispatch`) | Python/gRPC policy and native MPI coverage for app/workload sampling, CPU/GPU compatibility, runtime scaling, turnaround and IPDPS24 placement, adapted-limit and actual-duration wall-time policies, evaluation metrics, machine-size eligibility, and oversized-request truncation | | Backfill-window gRPC | 5 repeated checks | `run_backfill_window_grpc_test.sh` | Focused rerun of the gRPC streaming binary; one check targets the backfill window | | Single-coordinator gRPC | 1 | `test_grpc_single_coordinator.py` | Synchronized independent simulation servers | | Queue input schema | 1 binary | CTest or installed `test_queue_input` | Legacy queue names or numeric queue IDs, plus accepted and rejected replay/simulation runtime invariants | diff --git a/tests/run_append_job_tests.sh b/tests/run_append_job_tests.sh index 42c25f09..593bc727 100755 --- a/tests/run_append_job_tests.sh +++ b/tests/run_append_job_tests.sh @@ -8,7 +8,7 @@ # sitting in a preloaded m_data. See # docs/dev/OUTPUT_TRACE_BUFFERS.md for the design. # -# The in-process binary contains 20 focused append, batch, capacity, +# The in-process binary contains 23 focused append, batch, capacity, # advancement, accounting, and memory-pressure checks. When gRPC is built, # the runner also executes wire-level append, monitoring, and warm-start batch # checks against a real server. @@ -48,7 +48,7 @@ echo "" PASS=0 FAIL=0 -# --- Test: C++ API (test_append_job_api.cpp's 20 focused checks) --- +# --- Test: C++ API (test_append_job_api.cpp's 23 focused checks) --- echo "Testing: append_job_api (C++ level)" # Test binaries are installed under bin/tests/ (see CMakeLists.txt's diff --git a/tests/test_append_job_api.cpp b/tests/test_append_job_api.cpp index a1b45eb3..66a5af82 100644 --- a/tests/test_append_job_api.cpp +++ b/tests/test_append_job_api.cpp @@ -1076,6 +1076,54 @@ void test_append_with_known_actual_runtime() { std::cout << " PASSED" << std::endl; } +void test_fractional_streaming_limits_validate_stored_value() { + std::cout << "\n=== Fractional limits validate the stored value ===" + << std::endl; + Simulation sim(make_params()); + sim.get_trace().load_data(0); + + bool rejected = false; + try { + (void)sim.append_job(0.0, 1, kTestQueueInput, 10.25); + } catch (const std::invalid_argument &) { + rejected = true; + } + assert(rejected); + assert(sim.get_trace().data().empty()); + + const auto single = sim.append_job(0.0, 1, kTestQueueInput, 11, 10.25); + assert(sim.get_trace().job_at(single).get_limit_time() == 11); + assert(approx_equal(sim.get_trace().job_at(single).get_actual_run_time(), + 10.25)); + + const std::vector invalid_batch = { + {1.0, 1, kTestQueueInput, 21, 20.01}, + {2.0, 1, kTestQueueInput, 30.99}, + }; + rejected = false; + try { + (void)sim.append_jobs(invalid_batch); + } catch (const std::invalid_argument &) { + rejected = true; + } + assert(rejected); + assert(sim.get_trace().data().size() == 1); + + const std::vector batch = { + {1.0, 1, kTestQueueInput, 21, 20.01}, + {2.0, 1, kTestQueueInput, 31, 30.99}, + }; + const auto jobs = sim.append_jobs(batch); + assert(sim.get_trace().job_at(jobs[0]).get_limit_time() == 21); + assert(sim.get_trace().job_at(jobs[1]).get_limit_time() == 31); + assert(approx_equal(sim.get_trace().job_at(jobs[0]).get_actual_run_time(), + 20.01)); + assert(approx_equal(sim.get_trace().job_at(jobs[1]).get_actual_run_time(), + 30.99)); + + std::cout << " PASSED" << std::endl; +} + int main() { std::cout << "====================================" << std::endl; std::cout << "Append-Job Test Suite" << std::endl; @@ -1106,6 +1154,7 @@ int main() { test_resource_area_and_time_accounted_utilization(); test_prediction_horizon(); test_append_with_known_actual_runtime(); + test_fractional_streaming_limits_validate_stored_value(); std::cout << "\n====================================" << std::endl; std::cout << "ALL APPEND_JOB TESTS PASSED" << std::endl; diff --git a/tests/test_mpi_performance_dispatch.py b/tests/test_mpi_performance_dispatch.py index 91a0a37b..baa0e823 100644 --- a/tests/test_mpi_performance_dispatch.py +++ b/tests/test_mpi_performance_dispatch.py @@ -149,8 +149,8 @@ def main(): ) with actual_duration_output.open(newline="") as stream: for row in csv.DictReader(stream): - assert float(row["submitted_time_limit"]) == float( - row["actual_duration"] + assert float(row["submitted_time_limit"]) == math.ceil( + float(row["actual_duration"]) ) assert int(row["time_limit_doublings"]) == 0