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
7 changes: 7 additions & 0 deletions docs/CLIENT_SERVER_GUIDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ in-process via the streaming API:
| `RunUntilExclusiveRequest` | `Simulation::run_until_exclusive()` |
| `GetFCFSHeadShadowTimeRequest` | FCFS-head shadow time only: the earliest reserved start time, or `-1` with no waiting head |
| `GetBackfillWindowRequest` | One FCFS/EASY reservation snapshot: current capacity, shadow time, and projected releases |
| `GetPredictionHorizonRequest` | On-demand FCFS/EASY waiting-resource-time horizon after the shadow time |
| `GetStatisticsRequest`, `GetCurrentTimeRequest`, etc. | The monitoring/statistics methods |

For a replay-based warm start, set `InitRequest.sim_start_time` to a positive
Expand Down Expand Up @@ -86,6 +87,12 @@ a 60-node job predicted to end at time 100, and a 100-node FCFS head, the
response at time 0 has `shadow_time = 100` and releases `(50, 40)` and
`(100, 60)`.

Send `GetPredictionHorizonRequest` with a utilization factor in `[0, 1]` to
estimate how long the current waiting resource-time demand takes to drain after
the shadow time. Zero selects the fallback factor `1`. The server scans fields
already stored in supported FCFS queues only for this request; no prediction
state is maintained during normal scheduling.

Session initialization, completion, reuse, and shutdown are documented in
[Client/Server Setup](user-guide/grpc-setup.md#session-identity-and-completion).

Expand Down
4 changes: 2 additions & 2 deletions docs/api/PYTHON_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,8 +102,8 @@ Their scheduling semantics are documented in
| `get_available_nodes()` | Return free nodes. |
| `get_active_job_count()` | Return waiting jobs. |
| `get_fcfs_head_shadow_time()` | Return the FCFS-head reservation time, or `-1`. |
| `get_backfill_window()` | Return the current FCFS/EASY reservation snapshot. |
| `get_prediction_horizon(utilization)` | Estimate the Custom-FCFS/EASY waiting-queue drain time from the FCFS shadow time. |
| `get_backfill_window()` | Return the current FCFS/EASY reservation snapshot, including all running-job releases. |
| `get_prediction_horizon(utilization)` | On demand, estimate the supported FCFS/EASY waiting-queue drain time from the shadow time. |
| `get_statistics()` | Return a `Statistics` snapshot. |
| `write_simulated_trace()` | Write the configured job-schedule output. |
| `print_stats()` | Print summary statistics. |
Expand Down
26 changes: 16 additions & 10 deletions docs/api/STREAMING_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -265,8 +265,11 @@ The snapshot contains `current_time`, immediately `available_nodes`, the
same FCFS-head `shadow_time` (`-1` if the queue is empty), and chronologically
ordered resource-change events in `releases`. Each event gives the simulation
`time` at which capacity changes and the summed `nodes_released` then. Events
use time-limit estimates and extend through the reservation; simultaneous
releases are combined. This is an in-process API; it does not require gRPC.
use time-limit estimates and include every currently running job, even when no
FCFS head is waiting; simultaneous releases are combined. This is an
in-process API; it does not require gRPC. The gRPC request returns the same
full projection, which is useful when estimating when a prospective job could
acquire enough nodes.

```cpp
auto window = sim.get_backfill_window();
Expand All @@ -275,18 +278,20 @@ for (const auto& change : window.releases) {
}
```

**Estimate the Custom-FCFS waiting-queue prediction horizon:**
**Estimate the waiting-queue prediction horizon:**

```cpp
tdiff_t horizon = sim.get_prediction_horizon(utilization);
```

This method is available only for a simulation created with the Custom-FCFS
callback constructor and configured for EASY backfilling. Call it after the
current backfilling cycle completes. At that point, the waiting queue contains
only jobs that could not start, and the running set includes jobs dispatched by
the cycle. The method holds those sets fixed: future arrivals are excluded and
no additional waiting jobs are admitted during its forward replay.
This method is available for the standard and Custom FCFS implementations with
EASY backfilling. Call it after the current backfilling cycle completes. At
that point, the waiting queue contains only jobs that could not start, and the
running set includes jobs dispatched by the cycle. The method holds those sets
fixed: future arrivals are excluded and no additional waiting jobs are admitted
during its forward replay. Waiting resource-time is scanned from existing queue
records only when this method is called; normal scheduling maintains no extra
prediction state.

Queued demand is `A_Q = sum(requested_nodes * estimated_runtime)`. Starting at
the FCFS head's shadow time, the method integrates available nodes over each
Expand All @@ -300,7 +305,8 @@ The method returns the first completion-event offset at which accumulated
usable area covers `A_Q`; it does not interpolate within an intermediate
interval. If the final currently running job completes before the threshold
is reached, the remaining area is converted to time using
`utilization * total_nodes`. An empty queue returns zero.
`utilization * total_nodes`. An empty queue returns zero. The gRPC
`GetPredictionHorizonRequest` exposes the same calculation.

**Get scheduling statistics** (wait times, turnaround, utilization):
```cpp
Expand Down
40 changes: 40 additions & 0 deletions experimental/multi-cluster/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
# Performance-aware multi-cluster dispatch experiment

This study routes an arrival stream among independent DR_EVT systems by
predicted turnaround time. It is experimental code rather than a supported
general-purpose client.

For each arriving job, `grpc_performance_dispatch.py` finds the nearest row in
`performance_table.csv` using range-normalized Euclidean distance over node
count and time limit. It queries each server's backfill snapshot and queue
prediction horizon, then minimizes:

```text
predicted turnaround = predicted wait + time_limit / relative_performance
```

The performance table starts with `profile_id,num_nodes,time_limit`, followed
by one column per system. Each system value is a positive relative speed: `2.0`
means twice the baseline speed and half the processing time.

Run one controller and three simulation servers with:

```bash
mpirun -np 4 python3 python/grpc_mpi_launcher.py \
--server-binary "${CMAKE_INSTALL_PREFIX}/bin/dr_evt_server" \
--client-script experimental/multi-cluster/grpc_performance_dispatch.py -- \
--jobs python/examples/sample_trace.csv \
--performance-table experimental/multi-cluster/performance_table.csv \
--system-id system-1 --system-id system-2 --system-id system-3 \
--total-nodes 100 --prediction-utilization 1.0 \
--output dispatch-decisions.csv
```

This requires `mpi4py`, `grpcio`, `grpcio-tools`, and `protobuf`. The MPI
launcher uses rank 0 for the controller and one non-root rank per server.

Run the experiment's focused unit tests directly:

```bash
python3 experimental/multi-cluster/test_performance_dispatch.py
```
Loading
Loading