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
53 changes: 0 additions & 53 deletions .codex/skills/developer/SKILL.md

This file was deleted.

7 changes: 5 additions & 2 deletions docs/api/PYTHON_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,8 +160,11 @@ contains `time` and `nodes_released`.
`jobs_waiting`;
- `current_time`, `total_nodes`, `nodes_in_use`, and
`nodes_available`; and
- `resource_area`, `utilization`, `avg_wait_time`, `avg_turnaround_time`, and
`makespan`.
- `resource_area`, `utilization`, `avg_wait_time`, `avg_run_time`,
`avg_turnaround_time`, `avg_bounded_slowdown`, and `makespan`.

`avg_bounded_slowdown` is the mean of
`max(1, turnaround / max(run_time, 10 seconds))` over completed jobs.

For simulations constructed with Custom-FCFS callbacks, `resource_area` is
accumulated as `nodes_in_use * interval` between settled scheduling times.
Expand Down
4 changes: 4 additions & 0 deletions docs/api/STREAMING_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -381,6 +381,10 @@ is reached, the remaining area is converted to time using
`GetPredictionHorizonRequest` exposes the same calculation.

**Get scheduling statistics** (wait times, turnaround, utilization):

The returned statistics also include mean execution time and mean bounded
slowdown. Bounded slowdown is computed per completed job as
`max(1, turnaround / max(run_time, 10 seconds))`.
```cpp
Simulation::Statistics get_statistics() const;
```
Expand Down
69 changes: 46 additions & 23 deletions experimental/multi-cluster/README.md
Original file line number Diff line number Diff line change
@@ -1,16 +1,21 @@
# Online multi-cluster workload dispatch

`mpi_performance_dispatch` models one controller and one independent DR_EVT
simulation per system. It combines a chronological job stream with sampled
application workloads and dispatches every arrival using current simulated
queue state and predicted relative performance. Ground-truth relative
performance determines how long the selected job actually runs. The decision
is online: no routing plan is computed before the run.

The native implementation uses MPI and Ser20, but not gRPC or Python. Rank 0
owns the input tables, random-number generator, and dispatch policy. Each
other rank owns one `Simulation`; worker-rank order matches systems-table row
order.
`mpi_performance_dispatch` models one dispatcher and multiple cluster systems,
with one independent DR_EVT simulation per system. For each arriving workload,
the dispatcher selects a compatible system using the workload's predicted
relative performance and the systems' current simulated queue states.
Ground-truth relative performance determines how long the workload actually
runs on the selected system. Dispatch is online so that each decision accounts
for the latest estimated queue wait; the policy minimizes predicted turnaround,
which combines that wait with predicted run time.

The native implementation uses MPI, with Ser20 serialization for message
packing; it does not require gRPC or Python. Rank 0 owns the input tables,
random-number generator, and dispatch policy. Each other rank owns one
`Simulation`, and worker-rank order matches systems-table row order. The
corresponding Python implementation uses gRPC for dispatcher-worker
communication, while MPI starts the worker processes that host the gRPC
servers.

## Input model

Expand Down Expand Up @@ -124,17 +129,6 @@ needed. Quartz data is not required. In `merged.txt`, the unsuffixed `matrix`,
`tioga`, and `tuolumne` measurements represent their GPU modes; their `-cpu`
rows represent CPU modes.

Plot actual (x-axis) against predicted (y-axis) relative performance for every
mode, with a consistent color for each application:

```bash
source docs/venv/bin/activate
python3 experimental/multi-cluster/plot_relative_performance.py \
--ground-truth multi-cluster/ground_truth.csv \
--prediction multi-cluster/prediction.csv \
--output multi-cluster/relative-performance.png
```

Generate two interchangeable prediction baselines:

```bash
Expand All @@ -153,6 +147,22 @@ kernel names remain distinct positional samples. The application-average
baseline assigns the arithmetic mean ground-truth speedup for each
`(App, Ranks, execution mode)` group.

#### Optional plotting

Plot actual (x-axis) against predicted (y-axis) relative performance for every
configured machine/execution-mode column (for example, `dane`, `mammoth`,
`matrix-cpu`, or `matrix-gpu`), with a consistent color for each application.
Here, execution mode identifies the CPU or GPU implementation on a machine; it
does not refer to the simulator's `run_time_mode` option:

```bash
source docs/venv/bin/activate
python3 experimental/multi-cluster/plot_relative_performance.py \
--ground-truth multi-cluster/ground_truth.csv \
--prediction multi-cluster/prediction.csv \
--output multi-cluster/relative-performance.png
```

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.
Expand Down Expand Up @@ -248,7 +258,20 @@ Rank 0 writes one CSV row per job-stream row:
| `job_idx` | Worker-local DR_EVT job identifier. |

After draining all workers, rank 0 prints submitted/completed counts and
makespan per system to standard error.
makespan per system to standard error, followed by an `overall` line containing
these evaluation metrics:

| Metric | Definition |
|---|---|
| `average_turnaround_time` | Sum of actual submit-to-completion times divided by completed jobs. |
| `average_bounded_slowdown` | Mean of `max(1, turnaround / max(actual run time, 10 seconds))`. |
| `average_run_time` | Sum of ground-truth-scaled run times divided by dispatched jobs. |
| `average_speedup` | Sum of selected ground-truth workload speedups divided by dispatched jobs. |

Turnaround and bounded slowdown are weighted by each system's completed-job
count when combined. A mismatch between the completed and dispatched counts is
reported as an error rather than producing partial metrics. The Python/gRPC
implementation prints the same summary schema.

## Python/gRPC implementation

Expand Down
54 changes: 51 additions & 3 deletions experimental/multi-cluster/grpc_performance_dispatch.py
Original file line number Diff line number Diff line change
Expand Up @@ -532,14 +532,62 @@ def write_results(stream, decisions):
writer.writerows(decisions)


def write_summary(stream, statistics, system_ids):
"""Print one completion summary per simulation server."""
def evaluation_metrics(decisions, statistics):
"""Return completion-weighted metrics for the dispatched workload."""
job_count = len(decisions)
completed = sum(stats.jobs_completed for stats in statistics)
if completed != job_count:
raise ValueError(
f"completed job count {completed} does not match dispatched "
f"job count {job_count}"
)
if job_count == 0:
return {
"average_turnaround_time": 0.0,
"average_bounded_slowdown": 0.0,
"average_run_time": 0.0,
"average_speedup": 0.0,
}
return {
"average_turnaround_time": sum(
stats.avg_turnaround_time * stats.jobs_completed
for stats in statistics
)
/ completed,
"average_bounded_slowdown": sum(
stats.avg_bounded_slowdown * stats.jobs_completed
for stats in statistics
)
/ completed,
"average_run_time": sum(
decision["actual_duration"] for decision in decisions
)
/ job_count,
"average_speedup": sum(
decision["ground_truth_relative_performance"]
for decision in decisions
)
/ job_count,
}


def write_summary(stream, statistics, system_ids, decisions):
"""Print per-system completion data and overall evaluation metrics."""
for system_id, stats in zip(system_ids, statistics):
print(
f"{system_id}: submitted={stats.jobs_submitted} "
f"completed={stats.jobs_completed} makespan={stats.makespan:.6g}",
file=stream,
)
metrics = evaluation_metrics(decisions, statistics)
print(
f"overall: jobs={len(decisions)} "
f"average_turnaround_time={metrics['average_turnaround_time']:.8g} "
f"average_bounded_slowdown={metrics['average_bounded_slowdown']:.8g} "
f"average_run_time={metrics['average_run_time']:.8g} "
f"average_speedup={metrics['average_speedup']:.8g}",
file=stream,
)


def main():
Expand Down Expand Up @@ -583,7 +631,7 @@ def main():
write_results(stream, decisions)
else:
write_results(sys.stdout, decisions)
write_summary(sys.stderr, statistics, system_ids)
write_summary(sys.stderr, statistics, system_ids, decisions)
finally:
generated_dir.cleanup()

Expand Down
32 changes: 31 additions & 1 deletion experimental/multi-cluster/mpi_performance_dispatch.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -103,10 +103,13 @@ struct WindowMessage {
struct StatisticsMessage {
std::uint64_t jobs_submitted = 0;
std::uint64_t jobs_completed = 0;
double avg_turnaround_time = 0.0;
double avg_bounded_slowdown = 0.0;
double makespan = 0.0;

template <class Archive> void serialize(Archive &archive) {
archive(jobs_submitted, jobs_completed, makespan);
archive(jobs_submitted, jobs_completed, avg_turnaround_time,
avg_bounded_slowdown, makespan);
}
};

Expand Down Expand Up @@ -795,6 +798,8 @@ Response handle_request(dr_evt::Simulation &simulation,
const auto stats = simulation.get_statistics();
response.statistics = {static_cast<std::uint64_t>(stats.jobs_submitted),
static_cast<std::uint64_t>(stats.jobs_completed),
stats.avg_turnaround_time,
stats.avg_bounded_slowdown,
stats.makespan};
}
} catch (const std::exception &error) {
Expand Down Expand Up @@ -885,6 +890,9 @@ void controller(const Options &options, const std::vector<System> &systems,
"predicted_turnaround,job_idx\n";
*output << std::setprecision(17);

double total_run_time = 0.0;
double total_speedup = 0.0;

for (const auto &job : jobs) {
Job dispatch_job = job;
dispatch_job.num_nodes =
Expand All @@ -897,6 +905,8 @@ void controller(const Options &options, const std::vector<System> &systems,
const auto &workload = sample_workload(catalog, generator);
const Choice choice =
choose_system(dispatch_job, workload, systems, snapshots);
total_run_time += choice.actual_duration;
total_speedup += choice.ground_truth_relative_performance;

Request append;
append.operation = Operation::Append;
Expand Down Expand Up @@ -926,12 +936,32 @@ void controller(const Options &options, const std::vector<System> &systems,
Request finish;
finish.operation = Operation::Finish;
const auto responses = call_all(finish, worker_count);
std::uint64_t completed_jobs = 0;
double total_turnaround = 0.0;
double total_bounded_slowdown = 0.0;
for (std::size_t i = 0; i < responses.size(); ++i) {
const auto &stats = responses[i].statistics;
completed_jobs += stats.jobs_completed;
total_turnaround += stats.avg_turnaround_time * stats.jobs_completed;
total_bounded_slowdown +=
stats.avg_bounded_slowdown * stats.jobs_completed;
std::cerr << systems[i].id << ": submitted=" << stats.jobs_submitted
<< " completed=" << stats.jobs_completed
<< " makespan=" << stats.makespan << '\n';
}
if (completed_jobs != jobs.size())
throw std::runtime_error("completed job count does not match dispatched "
"job count");
const double denominator = static_cast<double>(jobs.size());
std::cerr << std::setprecision(8) << "overall: jobs=" << jobs.size()
<< " average_turnaround_time="
<< (jobs.empty() ? 0.0 : total_turnaround / denominator)
<< " average_bounded_slowdown="
<< (jobs.empty() ? 0.0 : total_bounded_slowdown / denominator)
<< " average_run_time="
<< (jobs.empty() ? 0.0 : total_run_time / denominator)
<< " average_speedup="
<< (jobs.empty() ? 0.0 : total_speedup / denominator) << '\n';
}

} // namespace
Expand Down
Loading
Loading