diff --git a/docs/CLIENT_SERVER_GUIDE.md b/docs/CLIENT_SERVER_GUIDE.md index fc5ea9d..d018d9e 100644 --- a/docs/CLIENT_SERVER_GUIDE.md +++ b/docs/CLIENT_SERVER_GUIDE.md @@ -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 @@ -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). diff --git a/docs/api/PYTHON_API.md b/docs/api/PYTHON_API.md index 099bab6..840acc5 100644 --- a/docs/api/PYTHON_API.md +++ b/docs/api/PYTHON_API.md @@ -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. | diff --git a/docs/api/STREAMING_API.md b/docs/api/STREAMING_API.md index 5646b1a..ece444b 100644 --- a/docs/api/STREAMING_API.md +++ b/docs/api/STREAMING_API.md @@ -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(); @@ -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 @@ -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 diff --git a/experimental/multi-cluster/README.md b/experimental/multi-cluster/README.md new file mode 100644 index 0000000..ffdc078 --- /dev/null +++ b/experimental/multi-cluster/README.md @@ -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 +``` diff --git a/experimental/multi-cluster/grpc_performance_dispatch.py b/experimental/multi-cluster/grpc_performance_dispatch.py new file mode 100644 index 0000000..cbc3c57 --- /dev/null +++ b/experimental/multi-cluster/grpc_performance_dispatch.py @@ -0,0 +1,305 @@ +#!/usr/bin/env python3 +"""Dispatch an arrival stream to independent DR_EVT systems by turnaround. + +The script is intended to be used as the rank-0 controller of +``grpc_mpi_launcher.py``. Each non-root MPI rank owns one DR_EVT server. For +every arrival, the controller finds the nearest performance-table profile, +queries every server's projected resource releases, and minimizes + + predicted turnaround = predicted wait + time_limit / relative_performance + +The selected server receives the job through AppendJobs and all servers stay +at the common arrival-time watermark through AdvanceTo. +""" + +import argparse +import csv +import math +import pathlib +import sys +from concurrent.futures import ThreadPoolExecutor + +# This experiment reuses the repository's general gRPC client helpers without +# making the experimental controller part of the Python package. +sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2] / "python")) + +from grpc_multi_server import (DEFAULT_QUEUE_INPUT, QUEUE_FIELD, ServerSession, + load_stubs) + + +REQUIRED_JOB_FIELDS = {"job_submit_time", "num_nodes", "time_limit"} +PROFILE_FIELDS = {"profile_id", "num_nodes", "time_limit"} + + +def read_arrivals(path): + """Read and validate a chronologically ordered simple-format trace.""" + with path.open(newline="") as stream: + reader = csv.DictReader(stream) + if not reader.fieldnames or not REQUIRED_JOB_FIELDS.issubset(reader.fieldnames): + raise ValueError( + f"{path} must have columns: {', '.join(sorted(REQUIRED_JOB_FIELDS))}" + ) + jobs = [] + for index, row in enumerate(reader): + job = { + "job_id": (row.get("job_id") or str(index)).strip(), + "submit_time": float(row["job_submit_time"]), + "num_nodes": int(row["num_nodes"]), + "queue": (row.get(QUEUE_FIELD) or DEFAULT_QUEUE_INPUT).strip(), + "limit_time": float(row["time_limit"]), + } + if job["num_nodes"] <= 0 or job["limit_time"] <= 0: + raise ValueError(f"{path}: job {job['job_id']} has non-positive size") + jobs.append(job) + if any(a["submit_time"] > b["submit_time"] for a, b in zip(jobs, jobs[1:])): + raise ValueError(f"{path}: jobs must be sorted by job_submit_time") + return jobs + + +def read_performance_table(path, system_ids): + """Read wide profiles whose system columns contain positive speedups.""" + required = PROFILE_FIELDS | set(system_ids) + with path.open(newline="") as stream: + reader = csv.DictReader(stream) + if not reader.fieldnames or not required.issubset(reader.fieldnames): + raise ValueError(f"{path} must have columns: {', '.join(sorted(required))}") + profiles = [] + for row in reader: + profile = { + "profile_id": row["profile_id"].strip(), + "num_nodes": float(row["num_nodes"]), + "time_limit": float(row["time_limit"]), + "performance": {name: float(row[name]) for name in system_ids}, + } + if not profile["profile_id"]: + raise ValueError(f"{path}: profile_id must not be empty") + if profile["num_nodes"] <= 0 or profile["time_limit"] <= 0: + raise ValueError(f"{path}: profile {profile['profile_id']} has non-positive size") + if any(value <= 0 or not math.isfinite(value) + for value in profile["performance"].values()): + raise ValueError( + f"{path}: profile {profile['profile_id']} has invalid performance" + ) + profiles.append(profile) + if not profiles: + raise ValueError(f"{path}: performance table is empty") + return profiles + + +def nearest_profile(job, profiles): + """Return nearest profile using range-normalized Euclidean distance.""" + node_values = [profile["num_nodes"] for profile in profiles] + time_values = [profile["time_limit"] for profile in profiles] + node_scale = max(node_values) - min(node_values) or 1.0 + time_scale = max(time_values) - min(time_values) or 1.0 + + def distance(profile): + node_delta = (job["num_nodes"] - profile["num_nodes"]) / node_scale + time_delta = (job["limit_time"] - profile["time_limit"]) / time_scale + return node_delta * node_delta + time_delta * time_delta + + return min(profiles, key=distance) + + +def estimate_release_wait(window, required_nodes): + """Return the first projected release with sufficient capacity.""" + available = window.available_nodes + if available >= required_nodes: + return 0.0 + for release in window.releases: + available += release.nodes_released + if available >= required_nodes: + return max(0.0, release.time - window.current_time) + return math.inf + + +def estimate_wait(window, required_nodes, runtime, prediction_horizon): + """Estimate wait from immediate EASY fit or queued resource-time.""" + has_waiting_head = window.shadow_time > window.current_time + can_start_now = window.available_nodes >= required_nodes + can_backfill_now = can_start_now and ( + not has_waiting_head + or window.current_time + runtime < window.shadow_time + ) + if can_backfill_now: + return 0.0 + if not has_waiting_head: + return estimate_release_wait(window, required_nodes) + return window.shadow_time - window.current_time + prediction_horizon + + +def choose_system(job, profile, system_ids, capacities, windows, horizons): + """Choose the feasible system with minimum predicted turnaround.""" + candidates = [] + for index, (system_id, capacity, window, horizon) in enumerate( + zip(system_ids, capacities, windows, horizons)): + performance = profile["performance"][system_id] + runtime = job["limit_time"] / performance + wait = estimate_wait(window, job["num_nodes"], runtime, horizon) + if job["num_nodes"] > capacity: + wait = math.inf + candidates.append({ + "index": index, + "system_id": system_id, + "relative_performance": performance, + "estimated_wait": wait, + "estimated_runtime": runtime, + "predicted_turnaround": wait + runtime, + }) + feasible = [candidate for candidate in candidates + if math.isfinite(candidate["predicted_turnaround"])] + if not feasible: + raise ValueError( + f"job {job['job_id']} ({job['num_nodes']} nodes) cannot fit any system" + ) + return min(feasible, key=lambda item: ( + item["predicted_turnaround"], item["estimated_wait"], item["index"])) + + +def call_all(executor, sessions, make_request): + futures = [executor.submit(session.call, make_request(session.messages)) + for session in sessions] + return [future.result() for future in futures] + + +def run_experiment(args, grpc, pb, service): + system_ids = args.system_id or [f"system-{index + 1}" + for index in range(len(args.server))] + if len(system_ids) != len(args.server) or len(set(system_ids)) != len(system_ids): + raise ValueError("--system-id must be unique and repeated once per --server") + capacities = args.system_nodes or [args.total_nodes] * len(args.server) + if len(capacities) != len(args.server) or any(value <= 0 for value in capacities): + raise ValueError("--system-nodes must be positive and repeated once per --server") + + jobs = read_arrivals(args.jobs) + profiles = read_performance_table(args.performance_table, system_ids) + sessions = [ServerSession(address, grpc, pb, service) for address in args.server] + decisions = [] + try: + with ThreadPoolExecutor(max_workers=len(sessions)) as executor: + init_futures = [] + for index, (session, capacity) in enumerate(zip(sessions, capacities)): + request = pb.ClientMessage(init=pb.InitRequest( + total_nodes=capacity, + trace_format="simple", + timestamp_format="epoch", + backfill_policy=args.backfill_policy, + priority_policy=args.priority_policy, + run_time_mode="limit", + infile=str(args.server_infile), + queue_impl=args.queue_impl, + session_name=f"{args.session_name}-{system_ids[index]}", + )) + init_futures.append(executor.submit(session.call, request)) + for future in init_futures: + future.result() + + for job in jobs: + call_all(executor, sessions, lambda messages: messages.ClientMessage( + advance_to=messages.AdvanceToRequest(target_time=job["submit_time"]))) + responses = call_all( + executor, sessions, lambda messages: messages.ClientMessage( + get_backfill_window=messages.GetBackfillWindowRequest())) + windows = [response.get_backfill_window for response in responses] + horizon_responses = call_all( + executor, sessions, lambda messages: messages.ClientMessage( + get_prediction_horizon=messages.GetPredictionHorizonRequest( + utilization=args.prediction_utilization))) + horizons = [response.get_prediction_horizon.horizon + for response in horizon_responses] + profile = nearest_profile(job, profiles) + choice = choose_system( + job, profile, system_ids, capacities, windows, horizons) + adjusted_runtime = choice["estimated_runtime"] + append = pb.AppendJobsRequest(requests=[pb.JobAppendData( + submit_time=job["submit_time"], num_nodes=job["num_nodes"], + queue=job["queue"], limit_time=adjusted_runtime)]) + response = sessions[choice["index"]].call( + pb.ClientMessage(append_jobs=append)) + sessions[choice["index"]].call(pb.ClientMessage( + advance_to=pb.AdvanceToRequest(target_time=job["submit_time"]))) + decisions.append({ + "job_id": job["job_id"], + "submit_time": job["submit_time"], + "profile_id": profile["profile_id"], + **{key: choice[key] for key in ( + "system_id", "relative_performance", "estimated_wait", + "estimated_runtime", "predicted_turnaround")}, + "job_idx": response.append_jobs.job_idx[0], + }) + + finish_responses = call_all( + executor, sessions, lambda messages: messages.ClientMessage( + finish_simulation=messages.FinishSimulationRequest())) + return decisions, [response.finish_simulation.statistics + for response in finish_responses], system_ids + finally: + for session in sessions: + session.close() + + +def write_results(stream, decisions): + fields = ("job_id", "submit_time", "profile_id", "system_id", + "relative_performance", "estimated_wait", "estimated_runtime", + "predicted_turnaround", "job_idx") + writer = csv.DictWriter(stream, fieldnames=fields) + writer.writeheader() + writer.writerows(decisions) + + +def write_summary(stream, statistics, system_ids): + 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) + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--server", action="append", required=True) + parser.add_argument("--system-id", action="append", + help="performance-table column; repeat in --server order") + parser.add_argument("--system-nodes", action="append", type=int, + help="node capacity; repeat in --server order") + parser.add_argument("--jobs", required=True, type=pathlib.Path) + parser.add_argument("--performance-table", required=True, type=pathlib.Path) + parser.add_argument("--server-infile", type=pathlib.Path) + parser.add_argument("--output", type=pathlib.Path, + help="decision CSV (default: stdout)") + parser.add_argument("--total-nodes", type=int, default=100) + parser.add_argument("--backfill-policy", default="easy") + parser.add_argument("--priority-policy", default="fcfs") + parser.add_argument("--queue-impl", default="circular") + parser.add_argument("--prediction-utilization", type=float, default=1.0, + help="usable-capacity factor for queue horizon (0..1)") + parser.add_argument("--session-name", default="performance-dispatch") + args = parser.parse_args() + if not 0.0 <= args.prediction_utilization <= 1.0: + parser.error("--prediction-utilization must be in [0, 1]") + if args.backfill_policy.lower() != "easy": + parser.error("performance dispatch requires --backfill-policy easy") + if args.priority_policy.lower() not in {"fcfs", "fcfs_alt"}: + parser.error("performance dispatch requires an FCFS priority policy") + args.server_infile = args.server_infile or args.jobs + + repo_root = pathlib.Path(__file__).resolve().parents[1] + grpc, pb, service, generated_dir = load_stubs(repo_root) + try: + decisions, statistics, system_ids = run_experiment( + args, grpc, pb, service) + if args.output: + with args.output.open("w", newline="") as stream: + write_results(stream, decisions) + else: + write_results(sys.stdout, decisions) + write_summary(sys.stderr, statistics, system_ids) + finally: + generated_dir.cleanup() + + +if __name__ == "__main__": + try: + main() + except (OSError, RuntimeError, ValueError) as error: + print(f"error: {error}", file=sys.stderr) + raise SystemExit(1) diff --git a/experimental/multi-cluster/performance_table.csv b/experimental/multi-cluster/performance_table.csv new file mode 100644 index 0000000..b504bf4 --- /dev/null +++ b/experimental/multi-cluster/performance_table.csv @@ -0,0 +1,4 @@ +profile_id,num_nodes,time_limit,system-1,system-2,system-3 +short-small,16,100,1.5,1.0,0.8 +medium,32,300,1.0,1.8,1.2 +large-long,64,900,0.7,1.1,2.0 diff --git a/experimental/multi-cluster/test_performance_dispatch.py b/experimental/multi-cluster/test_performance_dispatch.py new file mode 100644 index 0000000..5297cab --- /dev/null +++ b/experimental/multi-cluster/test_performance_dispatch.py @@ -0,0 +1,66 @@ +#!/usr/bin/env python3 +"""Unit tests for the multi-system performance dispatcher.""" + +import pathlib +import tempfile +import unittest +from types import SimpleNamespace + +from grpc_performance_dispatch import (choose_system, estimate_release_wait, + estimate_wait, nearest_profile, + read_performance_table) + + +def window(now, available, releases=(), shadow=-1): + return SimpleNamespace( + current_time=now, + available_nodes=available, + shadow_time=shadow, + releases=[SimpleNamespace(time=time, nodes_released=nodes) + for time, nodes in releases], + ) + + +class PerformanceDispatchTests(unittest.TestCase): + def test_normalized_nearest_profile(self): + profiles = [ + {"profile_id": "small", "num_nodes": 10, "time_limit": 100}, + {"profile_id": "large", "num_nodes": 90, "time_limit": 900}, + ] + job = {"num_nodes": 80, "limit_time": 800} + self.assertEqual(nearest_profile(job, profiles)["profile_id"], "large") + + def test_wait_uses_cumulative_releases(self): + snapshot = window(10, 2, ((15, 3), (21, 4))) + self.assertEqual(estimate_release_wait(snapshot, 5), 5) + self.assertEqual(estimate_release_wait(snapshot, 8), 11) + self.assertTrue(estimate_release_wait(snapshot, 10) == float("inf")) + + def test_wait_uses_shadow_and_queue_horizon(self): + snapshot = window(10, 2, shadow=30) + self.assertEqual(estimate_wait(snapshot, 8, 5, 40), 60) + immediately_backfillable = window(10, 8, shadow=30) + self.assertEqual(estimate_wait(immediately_backfillable, 8, 5, 40), 0) + + def test_turnaround_balances_wait_and_speed(self): + job = {"job_id": "j", "num_nodes": 8, "limit_time": 100} + profile = {"performance": {"fast": 4.0, "ready": 1.0}} + choice = choose_system( + job, profile, ["fast", "ready"], [10, 10], + [window(0, 2, ((80, 8),), shadow=80), window(0, 10)], + [0, 0]) + self.assertEqual(choice["system_id"], "ready") + self.assertEqual(choice["predicted_turnaround"], 100) + + def test_performance_table_requires_positive_values(self): + with tempfile.TemporaryDirectory() as directory: + path = pathlib.Path(directory) / "performance.csv" + path.write_text( + "profile_id,num_nodes,time_limit,a,b\n" + "cpu,4,20,2.0,0\n", encoding="utf-8") + with self.assertRaisesRegex(ValueError, "invalid performance"): + read_performance_table(path, ["a", "b"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/python/dr_evt_bindings.cpp b/python/dr_evt_bindings.cpp index 531d8cf..723da69 100644 --- a/python/dr_evt_bindings.cpp +++ b/python/dr_evt_bindings.cpp @@ -256,13 +256,13 @@ PYBIND11_MODULE(dr_evt, m) { .def("get_backfill_window", &Simulation::get_backfill_window, "Return a BackfillWindow snapshot for evaluating a backfill " - "candidate.") + "candidate, including all running-job releases.") .def("get_prediction_horizon", &Simulation::get_prediction_horizon, py::arg("utilization"), - "Estimate the Custom-FCFS waiting-queue drain time from the FCFS " - "shadow time. Requires EASY backfilling; future arrivals are " - "excluded.") + "Estimate the waiting-queue drain time from the FCFS shadow time. " + "Computed on demand for supported EASY schedulers; future arrivals " + "are excluded.") // Monitoring - Comprehensive statistics .def("get_statistics", &Simulation::get_statistics, diff --git a/src/proto/dr_evt_server.cpp b/src/proto/dr_evt_server.cpp index c479aa5..c40e1f5 100644 --- a/src/proto/dr_evt_server.cpp +++ b/src/proto/dr_evt_server.cpp @@ -390,6 +390,13 @@ class SimulationServiceImpl final : public SimulationService::Service { } break; } + case ClientMessage::kGetPredictionHorizon: { + require_init(sim); + resp.mutable_get_prediction_horizon()->set_horizon( + sim->get_prediction_horizon( + req.get_prediction_horizon().utilization())); + break; + } case ClientMessage::kGetStatistics: { require_init(sim); auto stats = sim->get_statistics(); diff --git a/src/proto/dr_evt_service.proto b/src/proto/dr_evt_service.proto index 5e38a8d..3deaea2 100644 --- a/src/proto/dr_evt_service.proto +++ b/src/proto/dr_evt_service.proto @@ -46,6 +46,7 @@ message ClientMessage { GetBackfillWindowRequest get_backfill_window = 16; GetCurrentUtilizationRequest get_current_utilization = 17; RunRequest run = 18; + GetPredictionHorizonRequest get_prediction_horizon = 19; } } @@ -160,6 +161,11 @@ message GetBackfillWindowRequest {} // Instantaneous nodes-in-use / total-nodes ratio. This is distinct from the // resource-area-based utilization in GetStatisticsResponse. message GetCurrentUtilizationRequest {} +// Estimate how long the currently waiting resource-time demand takes to drain +// after the FCFS-head shadow time. Computed only when queried. +message GetPredictionHorizonRequest { + double utilization = 1; +} message GetStatisticsRequest {} message GetTraceSizeRequest {} @@ -188,6 +194,7 @@ message ServerMessage { GetBackfillWindowResponse get_backfill_window = 17; GetCurrentUtilizationResponse get_current_utilization = 18; RunResponse run = 19; + GetPredictionHorizonResponse get_prediction_horizon = 20; } } @@ -243,6 +250,10 @@ message GetCurrentUtilizationResponse { double utilization = 1; } +message GetPredictionHorizonResponse { + double horizon = 1; +} + message GetActiveJobCountResponse { uint64 active_job_count = 1; } @@ -264,8 +275,8 @@ message GetBackfillWindowResponse { uint32 available_nodes = 2; // Earliest start of the FCFS queue head, or -1 when there is no head. double shadow_time = 3; - // Ordered release events in (current_time, shadow_time]. Empty when - // there is no FCFS head or it can start immediately. + // Ordered release events for all running jobs. May be nonempty when there + // is no FCFS head. repeated ResourceRelease releases = 4; } diff --git a/src/sim/block_wait_queue.hpp b/src/sim/block_wait_queue.hpp index dea7248..c1c983d 100644 --- a/src/sim/block_wait_queue.hpp +++ b/src/sim/block_wait_queue.hpp @@ -109,6 +109,12 @@ template class BlockWaitQueue { bool get_job_info(job_no_t job_id, tdiff_t &run_time, num_nodes_t &nodes) const; + /** + * Sum node-time for active jobs that have arrived by current_time. + * This scans existing queue records only when explicitly queried. + */ + tdiff_t waiting_resource_area(sim_time_t current_time) const; + /** * @brief Invoke a callable for every active job in FCFS order. * @details @@ -383,6 +389,25 @@ bool BlockWaitQueue::get_job_info(job_no_t job_id, tdiff_t &run_time, return false; } +template +tdiff_t BlockWaitQueue::waiting_resource_area( + sim_time_t current_time) const { + tdiff_t area = 0.0; + for (const auto &block_info : m_blocks) { + if (block_info.active_count == 0) { + continue; + } + const auto &seq = block_info.block.template get<0>(); + for (const auto &job : seq) { + if (job.submit_time <= current_time) { + area += + static_cast(job.nodes_requested) * job.run_time_estimate; + } + } + } + return area; +} + template template void BlockWaitQueue::for_each_active(Func &&func) const { diff --git a/src/sim/scheduler_base.cpp b/src/sim/scheduler_base.cpp index 0f05b03..8a6608a 100644 --- a/src/sim/scheduler_base.cpp +++ b/src/sim/scheduler_base.cpp @@ -18,12 +18,69 @@ #include "sim/scheduler_ljf.hpp" #include "sim/scheduler_sjf.hpp" #include +#include #include #include +#include +#include +#include #include namespace dr_evt { +tdiff_t SchedulerBase::prediction_horizon(const running_jobs_t &running_jobs, + sim_time_t current_time, + double utilization) const { + if (m_backfill_policy != BackfillPolicy::EASY) { + throw std::logic_error("prediction horizon requires EASY backfilling"); + } + if (!std::isfinite(utilization) || utilization < 0.0 || utilization > 1.0) { + throw std::invalid_argument("utilization must be finite and in [0, 1]"); + } + + const auto queued_area = waiting_resource_area(); + if (!queued_area.has_value()) { + throw std::logic_error( + "prediction horizon is unsupported by this scheduler"); + } + if (*queued_area <= 0.0) { + return 0.0; + } + if (m_total_nodes == 0) { + return std::numeric_limits::infinity(); + } + + const double effective_utilization = utilization == 0.0 ? 1.0 : utilization; + const sim_time_t shadow_time = + std::max(current_time, m_fcfs_reservation_time); + double available_nodes = static_cast(m_total_nodes); + std::map releases_by_time; + for (const auto &[job_id, job] : running_jobs) { + (void)job_id; + const sim_time_t end_time = job.start_time + job.run_time; + if (end_time > shadow_time) { + available_nodes -= static_cast(job.nodes); + releases_by_time[end_time] += job.nodes; + } + } + + tdiff_t usable_area = 0.0; + sim_time_t previous_time = shadow_time; + for (const auto &[release_time, nodes_released] : releases_by_time) { + usable_area += effective_utilization * available_nodes * + (release_time - previous_time); + if (usable_area >= *queued_area) { + return release_time - shadow_time; + } + available_nodes += static_cast(nodes_released); + previous_time = release_time; + } + + return previous_time - shadow_time + + (*queued_area - usable_area) / + (effective_utilization * static_cast(m_total_nodes)); +} + sim_time_t SchedulerBase::calculate_fcfs_reservation( num_nodes_t nodes_needed, num_nodes_t free_nodes, const running_jobs_t &running_jobs, sim_time_t current_time) { diff --git a/src/sim/scheduler_base.hpp b/src/sim/scheduler_base.hpp index 8021966..11a1668 100644 --- a/src/sim/scheduler_base.hpp +++ b/src/sim/scheduler_base.hpp @@ -15,6 +15,7 @@ #include "trace/trace.hpp" #include #include +#include #include namespace dr_evt { @@ -159,6 +160,24 @@ class SchedulerBase { return m_fcfs_reservation_time; } + /** + * @brief Sum requested-node time for jobs currently waiting. + * @details Implementations compute this on demand from their existing queue + * records. The default reports that prediction is unsupported, avoiding any + * storage or scheduling-path overhead for schedulers that do not opt in. + */ + virtual std::optional waiting_resource_area() const { + return std::nullopt; + } + + /** + * @brief Estimate waiting-queue drain time after the FCFS shadow time. + * @details Uses waiting_resource_area() and projected running-job releases. + * No prediction state is maintained between calls. + */ + tdiff_t prediction_horizon(const running_jobs_t &running_jobs, + sim_time_t current_time, double utilization) const; + protected: /** * Return total size of the wait queue (all jobs, ALL states). diff --git a/src/sim/scheduler_block_fcfs.hpp b/src/sim/scheduler_block_fcfs.hpp index 4b4f8ce..22b7f8a 100644 --- a/src/sim/scheduler_block_fcfs.hpp +++ b/src/sim/scheduler_block_fcfs.hpp @@ -121,6 +121,10 @@ class BlockQueueFCFSScheduler : public SchedulerBase { /** @copydoc SchedulerBase::has_eligible_jobs */ bool has_eligible_jobs() override { return active_job_count() > 0; } + std::optional waiting_resource_area() const override { + return m_wait_queue.waiting_resource_area(m_current_tracked_time); + } + protected: /** @copydoc SchedulerBase::wait_queue_size */ size_t wait_queue_size() const override { return m_wait_queue.size(); } diff --git a/src/sim/scheduler_circular_fcfs.hpp b/src/sim/scheduler_circular_fcfs.hpp index 2b90c56..129a2d0 100644 --- a/src/sim/scheduler_circular_fcfs.hpp +++ b/src/sim/scheduler_circular_fcfs.hpp @@ -155,6 +155,18 @@ class CircularBufferFCFSScheduler : public SchedulerBase { /** @copydoc SchedulerBase::has_eligible_jobs */ bool has_eligible_jobs() override { return active_job_count() > 0; } + std::optional waiting_resource_area() const override { + tdiff_t area = 0.0; + for (size_t i = 0; i < m_eligible_end_idx; ++i) { + const auto &job = m_wait_queue[i]; + if (!job.removed) { + area += + static_cast(job.nodes_requested) * job.run_time_estimate; + } + } + return area; + } + protected: /** @copydoc SchedulerBase::wait_queue_size */ size_t wait_queue_size() const override { return m_wait_queue.size(); } diff --git a/src/sim/scheduler_fcfs.hpp b/src/sim/scheduler_fcfs.hpp index 8a9670a..7ce8abc 100644 --- a/src/sim/scheduler_fcfs.hpp +++ b/src/sim/scheduler_fcfs.hpp @@ -132,6 +132,18 @@ class FCFSScheduler : public SchedulerBase { /** @copydoc SchedulerBase::has_eligible_jobs */ bool has_eligible_jobs() override { return active_job_count() > 0; } + std::optional waiting_resource_area() const override { + tdiff_t area = 0.0; + for (size_t i = 0; i < m_eligible_end_idx; ++i) { + const auto &job = m_wait_queue[i]; + if (!job.removed) { + area += + static_cast(job.nodes_requested) * job.run_time_estimate; + } + } + return area; + } + protected: /** * Return total size of wait queue (ALL jobs, all states). diff --git a/src/sim/scheduler_fcfs_alt.cpp b/src/sim/scheduler_fcfs_alt.cpp index 065ee93..2fb18d8 100644 --- a/src/sim/scheduler_fcfs_alt.cpp +++ b/src/sim/scheduler_fcfs_alt.cpp @@ -14,6 +14,17 @@ namespace dr_evt { +std::optional FCFSAltScheduler::waiting_resource_area() const { + tdiff_t area = 0.0; + for (const auto &[submit_time, job] : m_wait_queue) { + (void)submit_time; + if (m_eligible_jobs.contains(job.job_id)) { + area += static_cast(job.nodes) * job.run_time; + } + } + return area; +} + FCFSAltScheduler::FCFSAltScheduler(num_nodes_t total_nodes, BackfillPolicy backfill_policy) : SchedulerBase(total_nodes, backfill_policy), m_current_tracked_time(0.0) { @@ -98,14 +109,14 @@ FCFSAltScheduler::schedule(num_nodes_t free_nodes, // FCFS head blocked - try backfilling // Calculate FCFS reservation time - sim_time_t reservation_time = calculate_fcfs_reservation( + m_fcfs_reservation_time = calculate_fcfs_reservation( head_nodes, free_nodes, running_jobs, current_time); - if (reservation_time <= current_time) { + if (m_fcfs_reservation_time <= current_time) { return {}; // No valid reservation window } - tdiff_t backfill_window = reservation_time - current_time; + tdiff_t backfill_window = m_fcfs_reservation_time - current_time; // Scan in FCFS order (earliest submit_time first) for backfill candidates for (auto it = m_wait_queue.begin(); it != m_wait_queue.end(); ++it) { diff --git a/src/sim/scheduler_fcfs_alt.hpp b/src/sim/scheduler_fcfs_alt.hpp index 7aec596..cc2c8e9 100644 --- a/src/sim/scheduler_fcfs_alt.hpp +++ b/src/sim/scheduler_fcfs_alt.hpp @@ -80,6 +80,8 @@ class FCFSAltScheduler : public SchedulerBase { /** @copydoc SchedulerBase::has_eligible_jobs */ bool has_eligible_jobs() override { return !m_eligible_jobs.empty(); } + std::optional waiting_resource_area() const override; + protected: /** @copydoc SchedulerBase::wait_queue_size */ size_t wait_queue_size() const override { return m_wait_queue.size(); } diff --git a/src/sim/scheduler_fcfs_custom.cpp b/src/sim/scheduler_fcfs_custom.cpp index 65537dd..b493967 100644 --- a/src/sim/scheduler_fcfs_custom.cpp +++ b/src/sim/scheduler_fcfs_custom.cpp @@ -123,62 +123,15 @@ double CustomFCFSScheduler::utilization_through(sim_time_t through_time) const { (static_cast(m_total_nodes) * duration); } -tdiff_t -CustomFCFSScheduler::prediction_horizon(const running_jobs_t &running_jobs, - sim_time_t current_time, - double utilization) const { - if (m_backfill_policy != BackfillPolicy::EASY) { - throw std::logic_error( - "prediction horizon requires Custom FCFS with EASY backfilling"); - } - if (!std::isfinite(utilization) || utilization < 0.0 || utilization > 1.0) { - throw std::invalid_argument("utilization must be finite and in [0, 1]"); - } - const double effective_utilization = utilization == 0.0 ? 1.0 : utilization; - - tdiff_t queued_area = 0.0; +std::optional CustomFCFSScheduler::waiting_resource_area() const { + tdiff_t area = 0.0; for (size_t i = 0; i < m_eligible_end_idx; ++i) { const auto &job = m_wait_queue[i]; if (!job.removed) { - queued_area += - static_cast(job.nodes_requested) * job.run_time_estimate; - } - } - if (queued_area <= 0.0) { - return 0.0; - } - if (m_total_nodes == 0) { - return std::numeric_limits::infinity(); - } - - const sim_time_t shadow_time = - std::max(current_time, m_fcfs_reservation_time); - double available_nodes = static_cast(m_total_nodes); - std::map releases_by_time; - for (const auto &[job_id, job] : running_jobs) { - (void)job_id; - const sim_time_t end_time = job.start_time + job.run_time; - if (end_time > shadow_time) { - available_nodes -= static_cast(job.nodes); - releases_by_time[end_time] += job.nodes; - } - } - - tdiff_t usable_area = 0.0; - sim_time_t previous_time = shadow_time; - for (const auto &[release_time, nodes_released] : releases_by_time) { - usable_area += effective_utilization * available_nodes * - (release_time - previous_time); - if (usable_area >= queued_area) { - return release_time - shadow_time; + area += static_cast(job.nodes_requested) * job.run_time_estimate; } - available_nodes += static_cast(nodes_released); - previous_time = release_time; } - - return previous_time - shadow_time + - (queued_area - usable_area) / - (effective_utilization * static_cast(m_total_nodes)); + return area; } std::optional CustomFCFSScheduler::select_backfill_candidate( diff --git a/src/sim/scheduler_fcfs_custom.hpp b/src/sim/scheduler_fcfs_custom.hpp index 3c3f1fb..71dc280 100644 --- a/src/sim/scheduler_fcfs_custom.hpp +++ b/src/sim/scheduler_fcfs_custom.hpp @@ -104,9 +104,8 @@ class CustomFCFSScheduler : public SchedulerBase { */ double utilization_through(sim_time_t through_time) const; - /** Estimate the waiting-queue horizon for the settled Custom-FCFS state. */ - tdiff_t prediction_horizon(const running_jobs_t &running_jobs, - sim_time_t current_time, double utilization) const; + /** Compute queued node-time on demand from existing queue entries. */ + std::optional waiting_resource_area() const override; /** * Construct the FCFS scheduling core for a subclass-owned selection policy. diff --git a/src/sim/sim.cpp b/src/sim/sim.cpp index d97090b..f09e7a0 100644 --- a/src/sim/sim.cpp +++ b/src/sim/sim.cpp @@ -1262,19 +1262,13 @@ BasicSimulation::get_backfill_window() const { Backfill_Window window{ m_current_time, get_available_nodes(), get_fcfs_head_shadow_time(), {}}; - // A head that can start now has no future window to describe. Likewise, - // without a waiting head there is no EASY reservation. - if (window.shadow_time <= m_current_time) { - return window; - } - // m_running_jobs stores the actual start time. The scheduler reserves // against each job's limit time, so this deliberately does the same. std::map releases_by_time; for (const auto &[job_idx, job] : m_running_jobs) { (void)job_idx; const sim_time_t end_time = job.start_time + job.run_time; - if (end_time > m_current_time && end_time <= window.shadow_time) { + if (end_time > m_current_time) { releases_by_time[end_time] += job.nodes; } } @@ -1289,12 +1283,8 @@ BasicSimulation::get_backfill_window() const { template tdiff_t BasicSimulation::get_prediction_horizon(double utilization) const { - if (m_custom_scheduler == nullptr) { - throw std::logic_error( - "prediction horizon is available only with Custom FCFS"); - } - return m_custom_scheduler->prediction_horizon(m_running_jobs, m_current_time, - utilization); + return m_scheduler->prediction_horizon(m_running_jobs, m_current_time, + utilization); } template diff --git a/src/sim/sim.hpp b/src/sim/sim.hpp index 488cd5f..9d0ef15 100644 --- a/src/sim/sim.hpp +++ b/src/sim/sim.hpp @@ -375,9 +375,9 @@ template class BasicSimulation { * @details Returned by get_backfill_window() for a caller evaluating * whether a candidate can backfill without delaying the FCFS queue head. * `current_time` and `available_nodes` describe capacity immediately; - * `releases` then describes projected capacity increases up to the head's - * `shadow_time`. Release events use the time-limit estimates used by - * SchedulerBase::calculate_fcfs_reservation, rather than actual runtimes, + * `releases` then describes every projected capacity increase from the + * currently running jobs. Release events use the time-limit estimates used + * by SchedulerBase::calculate_fcfs_reservation, rather than actual runtimes, * so the projection and reservation agree. `shadow_time` is -1 when no * FCFS head is waiting. */ @@ -398,8 +398,7 @@ template class BasicSimulation { num_nodes_t available_nodes; ///< Nodes free immediately at current_time. sim_time_t shadow_time; ///< Reserved FCFS-head start time, or -1 if no head waits. - std::vector - releases; ///< Capacity increases through shadow_time. + std::vector releases; ///< Projected capacity increases. }; /** @@ -407,7 +406,8 @@ template class BasicSimulation { * @details * The result is a snapshot: available nodes and release times reflect * current scheduler state, while release times are based on the same - * time-limit estimates used for the FCFS reservation. + * time-limit estimates used for the FCFS reservation. Releases include + * every currently running job, including when no FCFS head waits. * @return Backfill_Window value for the current simulation time. */ Backfill_Window get_backfill_window() const; @@ -439,8 +439,8 @@ template class BasicSimulation { * * @pre The caller has completed scheduling at the current time, normally by * calling advance_to(get_current_time()) after adding jobs at that time. - * @pre This simulation was constructed with the CustomFCFSScheduler callback - * constructor and uses EASY backfilling. + * @pre The selected scheduler supports on-demand waiting-resource-area + * queries and uses EASY backfilling. * @param[in] utilization Expected system-utilization factor in (0, 1], or * zero to use the fallback factor 1. * @return Horizon duration from the shadow time. Returns zero for an empty @@ -448,8 +448,8 @@ template class BasicSimulation { * resources. * @throws std::invalid_argument if utilization is non-finite, negative, or * greater than 1. - * @throws std::logic_error if this simulation does not use Custom FCFS with - * EASY backfilling. + * @throws std::logic_error if the scheduler does not support prediction or + * does not use EASY backfilling. */ tdiff_t get_prediction_horizon(double utilization) const; diff --git a/tests/run_warm_start_validation_tests.sh b/tests/run_warm_start_validation_tests.sh index 6e8c974..7f61803 100755 --- a/tests/run_warm_start_validation_tests.sh +++ b/tests/run_warm_start_validation_tests.sh @@ -5,6 +5,7 @@ set -u if [ "$#" -ne 1 ]; then echo "Usage: $0 " >&2 + echo "Example: $0 ${CMAKE_INSTALL_PREFIX}/bin/simulator" >&2 exit 2 fi diff --git a/tests/test_append_job_api.cpp b/tests/test_append_job_api.cpp index 2c05dbc..143a0d9 100644 --- a/tests/test_append_job_api.cpp +++ b/tests/test_append_job_api.cpp @@ -948,16 +948,34 @@ void test_prediction_horizon() { empty.get_trace().load_data(0); assert(approx_equal(empty.get_prediction_horizon(0.5), 0.0)); - auto standard_params = make_params(); - Simulation standard(standard_params); - standard.get_trace().load_data(0); - bool rejected_for_standard_scheduler = false; + // Every standard FCFS queue computes the same demand by scanning only its + // existing eligible entries when this query is made. + for (const auto queue_impl : + {QueueImplementation::CIRCULAR, QueueImplementation::DEQUE, + QueueImplementation::MULTIMAP, QueueImplementation::BLOCK}) { + auto standard_params = make_params(); + standard_params.m_queue_impl = queue_impl; + Simulation standard(standard_params); + standard.get_trace().load_data(0); + standard.append_job(0.0, 100, kTestQueueInput, 40.0); + standard.advance_to(0.0); + standard.append_job(0.0, 50, kTestQueueInput, 8.0); + standard.advance_to(0.0); + assert(approx_equal(standard.get_fcfs_head_shadow_time(), 40.0)); + assert(approx_equal(standard.get_prediction_horizon(0.5), 8.0)); + } + + auto unsupported_params = make_params(); + unsupported_params.m_priority_policy = PriorityPolicy::SJF; + Simulation unsupported(unsupported_params); + unsupported.get_trace().load_data(0); + bool rejected_for_unsupported_scheduler = false; try { - (void)standard.get_prediction_horizon(0.5); + (void)unsupported.get_prediction_horizon(0.5); } catch (const std::logic_error &) { - rejected_for_standard_scheduler = true; + rejected_for_unsupported_scheduler = true; } - assert(rejected_for_standard_scheduler); + assert(rejected_for_unsupported_scheduler); auto no_backfill_params = make_params(); no_backfill_params.m_backfill_policy = BackfillPolicy::NONE; diff --git a/tests/test_grpc_streaming_api.cpp b/tests/test_grpc_streaming_api.cpp index 3f92405..fe1434b 100644 --- a/tests/test_grpc_streaming_api.cpp +++ b/tests/test_grpc_streaming_api.cpp @@ -341,6 +341,16 @@ bool test_backfill_window(const std::string &server_address, client.finish(); return false; } + + ClientMessage horizon_query; + horizon_query.mutable_get_prediction_horizon()->set_utilization(1.0); + const auto horizon_response = client.call(horizon_query); + if (std::fabs(horizon_response.get_prediction_horizon().horizon() - 200.0) > + 1e-12) { + std::cerr << " FAIL: expected prediction horizon 200\n"; + client.finish(); + return false; + } } catch (const std::exception &e) { std::cerr << " FAIL: " << e.what() << "\n"; client.finish(); diff --git a/tests/test_python_api.py b/tests/test_python_api.py index b772138..75e1254 100755 --- a/tests/test_python_api.py +++ b/tests/test_python_api.py @@ -321,6 +321,20 @@ def test_backfill_window_api(result): # No running-job completion remains after the shadow event, so the # 1000 node-seconds are drained at U * total_nodes = 50 nodes. assert abs(sim.get_prediction_horizon(0.5) - 20.0) < 1e-12 + # Once the waiting head starts, the snapshot continues to expose the + # projected releases of running work. + sim.advance_to(100.0) + full_window = sim.get_backfill_window() + assert full_window.shadow_time == -1.0 + assert [(release.time, release.nodes_released) + for release in full_window.releases] == [(110.0, 100)] + + standard = dr_evt.Simulation(params) + standard.append_job(0.0, 100, QUEUE_INPUT, 40.0) + standard.advance_to(0.0) + standard.append_job(0.0, 50, QUEUE_INPUT, 8.0) + standard.advance_to(0.0) + assert standard.get_prediction_horizon(0.5) == 8.0 result.record_pass("Backfill window snapshot") except Exception as e: result.record_fail("Backfill window API", str(e))