2026-05-11 16:48:52 +01:00
|
|
|
import argparse
|
|
|
|
|
import asyncio
|
|
|
|
|
import logging
|
|
|
|
|
import os
|
|
|
|
|
import sys
|
2026-05-11 17:08:44 +01:00
|
|
|
import time
|
2026-05-11 16:48:52 +01:00
|
|
|
|
|
|
|
|
from vastai import Serverless
|
|
|
|
|
|
|
|
|
|
logging.basicConfig(
|
2026-05-11 17:08:44 +01:00
|
|
|
level=logging.INFO,
|
2026-05-11 16:48:52 +01:00
|
|
|
format="%(asctime)s[%(levelname)-5s] %(message)s",
|
|
|
|
|
datefmt="%Y-%m-%d %H:%M:%S",
|
|
|
|
|
)
|
|
|
|
|
log = logging.getLogger(__file__)
|
|
|
|
|
|
|
|
|
|
ENDPOINT_NAME = "null-prod"
|
|
|
|
|
|
|
|
|
|
|
2026-05-11 17:08:44 +01:00
|
|
|
async def reserve(
|
|
|
|
|
client: Serverless,
|
|
|
|
|
*,
|
|
|
|
|
endpoint_name: str,
|
|
|
|
|
duration: float,
|
|
|
|
|
label: str = "reservation",
|
|
|
|
|
) -> dict:
|
2026-05-11 16:48:52 +01:00
|
|
|
"""Hold a Vast worker open for `duration` seconds (or until we disconnect).
|
|
|
|
|
|
2026-05-11 17:08:44 +01:00
|
|
|
The worker counts itself busy for the lifetime of this call. Returning
|
|
|
|
|
here means the reservation has ended — either /release was called on
|
|
|
|
|
the worker's internal control port, or the duration cap fired, or the
|
|
|
|
|
HTTP request was cancelled.
|
2026-05-11 16:48:52 +01:00
|
|
|
"""
|
|
|
|
|
endpoint = await client.get_endpoint(name=endpoint_name)
|
|
|
|
|
payload = {"duration": duration}
|
2026-05-11 17:08:44 +01:00
|
|
|
start = time.monotonic()
|
|
|
|
|
log.info("[%s] POST /reserve duration=%ss", label, duration)
|
|
|
|
|
try:
|
2026-05-11 17:47:19 +01:00
|
|
|
resp = await endpoint.request("/reserve", payload, cost=1)
|
2026-05-11 17:08:44 +01:00
|
|
|
elapsed = time.monotonic() - start
|
|
|
|
|
log.info("[%s] returned after %.1fs: %s", label, elapsed, resp.get("response"))
|
|
|
|
|
return resp["response"]
|
|
|
|
|
except asyncio.CancelledError:
|
|
|
|
|
elapsed = time.monotonic() - start
|
|
|
|
|
log.info("[%s] cancelled after %.1fs (HTTP connection dropped)", label, elapsed)
|
|
|
|
|
raise
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def run_demo(
|
|
|
|
|
client: Serverless,
|
|
|
|
|
*,
|
|
|
|
|
endpoint_name: str,
|
|
|
|
|
interval: float,
|
|
|
|
|
) -> None:
|
2026-05-11 18:00:46 +01:00
|
|
|
"""Trapezoidal load: ramp up three reservations, then let them scale down.
|
|
|
|
|
|
|
|
|
|
Start three reservations spaced `interval` seconds apart, each with a
|
|
|
|
|
duration equal to 3 * interval. The staggered starts and identical
|
|
|
|
|
durations mean they end one at a time, also `interval` apart, so the
|
|
|
|
|
load curve ramps up over 2*interval, plateaus at 3 for `interval`, and
|
|
|
|
|
ramps down over 2*interval. Each reservation ends via its duration cap
|
|
|
|
|
(a 200 success, not a 499 cancellation).
|
2026-05-11 17:08:44 +01:00
|
|
|
"""
|
2026-05-11 18:00:46 +01:00
|
|
|
hold = interval * 3
|
2026-05-11 17:08:44 +01:00
|
|
|
tasks: list[asyncio.Task] = []
|
|
|
|
|
for i in range(1, 4):
|
|
|
|
|
label = f"res-{i}"
|
2026-05-11 18:00:46 +01:00
|
|
|
log.info(
|
|
|
|
|
"[%s] starting (auto-release after %.0fs)", label, hold
|
|
|
|
|
)
|
2026-05-11 17:08:44 +01:00
|
|
|
task = asyncio.create_task(
|
|
|
|
|
reserve(
|
|
|
|
|
client,
|
|
|
|
|
endpoint_name=endpoint_name,
|
2026-05-11 18:00:46 +01:00
|
|
|
duration=hold,
|
2026-05-11 17:08:44 +01:00
|
|
|
label=label,
|
|
|
|
|
),
|
|
|
|
|
name=label,
|
|
|
|
|
)
|
|
|
|
|
tasks.append(task)
|
|
|
|
|
if i < 3:
|
2026-05-11 18:00:46 +01:00
|
|
|
log.info("Waiting %.0fs before next reservation...", interval)
|
2026-05-11 17:08:44 +01:00
|
|
|
await asyncio.sleep(interval)
|
|
|
|
|
|
|
|
|
|
log.info(
|
2026-05-11 18:00:46 +01:00
|
|
|
"All 3 reservations in flight; they will scale down %.0fs apart, "
|
|
|
|
|
"starting in %.0fs",
|
2026-05-11 17:08:44 +01:00
|
|
|
interval,
|
2026-05-11 18:00:46 +01:00
|
|
|
hold - 2 * interval,
|
2026-05-11 17:08:44 +01:00
|
|
|
)
|
|
|
|
|
results = await asyncio.gather(*tasks, return_exceptions=True)
|
|
|
|
|
for task, result in zip(tasks, results):
|
|
|
|
|
log.info("[%s] final: %r", task.get_name(), result)
|
2026-05-11 16:48:52 +01:00
|
|
|
|
|
|
|
|
|
|
|
|
|
def build_arg_parser() -> argparse.ArgumentParser:
|
|
|
|
|
p = argparse.ArgumentParser(description="Vast Null PyWorker demo client")
|
|
|
|
|
p.add_argument(
|
|
|
|
|
"--endpoint",
|
|
|
|
|
default=os.environ.get("VAST_ENDPOINT", ENDPOINT_NAME),
|
|
|
|
|
help=f"Vast endpoint name (default: {ENDPOINT_NAME})",
|
|
|
|
|
)
|
|
|
|
|
p.add_argument(
|
|
|
|
|
"--duration",
|
|
|
|
|
type=float,
|
2026-05-11 17:08:44 +01:00
|
|
|
default=180.0,
|
|
|
|
|
help="Seconds to hold each worker busy (default: 180)",
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
modes = p.add_mutually_exclusive_group(required=False)
|
|
|
|
|
modes.add_argument(
|
|
|
|
|
"--reserve",
|
|
|
|
|
action="store_true",
|
|
|
|
|
help="Make a single /reserve call (default if no mode given)",
|
|
|
|
|
)
|
|
|
|
|
modes.add_argument(
|
|
|
|
|
"--demo",
|
|
|
|
|
action="store_true",
|
|
|
|
|
help="Run the staggered 3-reservation demo, cancelling one mid-way",
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
p.add_argument(
|
|
|
|
|
"--interval",
|
|
|
|
|
type=float,
|
|
|
|
|
default=30.0,
|
|
|
|
|
help="Demo mode: seconds between reservation steps (default: 30)",
|
2026-05-11 16:48:52 +01:00
|
|
|
)
|
|
|
|
|
return p
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def main_async():
|
|
|
|
|
args = build_arg_parser().parse_args()
|
|
|
|
|
|
|
|
|
|
print("=" * 60)
|
2026-05-11 17:08:44 +01:00
|
|
|
print(f"Endpoint: {args.endpoint}")
|
2026-05-11 16:48:52 +01:00
|
|
|
print("=" * 60)
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
async with Serverless() as client:
|
2026-05-11 17:08:44 +01:00
|
|
|
if args.demo:
|
|
|
|
|
await run_demo(
|
|
|
|
|
client,
|
|
|
|
|
endpoint_name=args.endpoint,
|
|
|
|
|
interval=args.interval,
|
|
|
|
|
)
|
|
|
|
|
else:
|
|
|
|
|
response = await reserve(
|
|
|
|
|
client,
|
|
|
|
|
endpoint_name=args.endpoint,
|
|
|
|
|
duration=args.duration,
|
|
|
|
|
label="reservation",
|
|
|
|
|
)
|
|
|
|
|
print(f"Reservation result: {response}")
|
|
|
|
|
except KeyboardInterrupt:
|
|
|
|
|
log.info("Interrupted; dropping any in-flight reservations")
|
2026-05-11 16:48:52 +01:00
|
|
|
except Exception as e:
|
2026-05-11 17:08:44 +01:00
|
|
|
log.error("Error: %s", e, exc_info=True)
|
2026-05-11 16:48:52 +01:00
|
|
|
sys.exit(1)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if __name__ == "__main__":
|
|
|
|
|
asyncio.run(main_async())
|