Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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 |
Expand Down
5 changes: 5 additions & 0 deletions .github/workflows/quick-test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 "========================================"
4 changes: 2 additions & 2 deletions .github/workflows/tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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 }}"
Expand Down
2 changes: 1 addition & 1 deletion docs/api/PYTHON_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down
4 changes: 3 additions & 1 deletion docs/api/STREAMING_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand All @@ -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)),
Expand Down
2 changes: 1 addition & 1 deletion docs/user-guide/trace-formats.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
33 changes: 22 additions & 11 deletions experimental/multi-cluster/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
```

Expand All @@ -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:
Expand All @@ -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
Expand All @@ -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. |

Expand Down Expand Up @@ -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
```

Expand Down
15 changes: 9 additions & 6 deletions experimental/multi-cluster/grpc_performance_dispatch.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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")
Expand Down
19 changes: 11 additions & 8 deletions experimental/multi-cluster/mpi_performance_dispatch.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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
};

Expand Down Expand Up @@ -171,7 +171,7 @@ struct Options {
double prediction_utilization = 1.0;
double max_time_limit = std::numeric_limits<double>::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 = {}) {
Expand All @@ -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);
}
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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);
Expand Down
18 changes: 15 additions & 3 deletions experimental/multi-cluster/plot_relative_performance.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import csv
import math
import pathlib
import sys


IDENTITY_COLUMNS = {"App", "Args", "Ranks", "app", "args", "ranks"}
Expand Down Expand Up @@ -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
):
Expand All @@ -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 (
Expand All @@ -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


Expand Down Expand Up @@ -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)


Expand Down
Loading
Loading