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
1 change: 1 addition & 0 deletions docs/dev-fleet-warm-slots.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ The contract is:
- Dirty/unavailable/recovery-required source returns cold_fallback_required so the controller can use its existing clean exact-SHA lane.
- A successful task build that consumed the shared warm lineage marks it warm_ready=false. The slot must be warmed back to main before it can advertise another task base.
- Reservations expire after a bounded lease interval. An abandoned task can release its exact lease explicitly; mismatched task/lease IDs fail closed.
- Recursive cache-size measurement is diagnostic work, not part of ordinary foreground execution. Normal `warm` and `task-run` calls do not walk the cache tree before/after native work. The physical benchmark passes `--measure-disk` explicitly when it needs cache-growth evidence, so large DerivedData trees cannot make routine telemetry a foreground latency tax.

The helper writes inflight.json before launching native work. A pipe launch guard keeps the child from executing the native command until its process group is durably recorded in both the in-flight record and visible lease. SIGINT/SIGTERM is forwarded to that group. If the helper dies unexpectedly, the guarded child exits before native exec or recover uses the exact recorded run/group identity. Recovery quarantines the lineage; diagnostic request records also carry process-start identity so PID reuse cannot keep a slot falsely busy.

Expand Down
2 changes: 2 additions & 0 deletions scripts/benchmark-dev-fleet-warm-slots.py
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,7 @@ def warm(
[
"warm", "--machine-state", str(state), "--slot", slot,
"--checkout", str(checkout), "--target", target,
"--measure-disk",
*command_tail(command),
],
env=env,
Expand All @@ -199,6 +200,7 @@ def task(
argv = [
"task-run", "--machine-state", str(state), "--slot", slot,
"--checkout", str(checkout), "--target", target, "--task-id", task_id,
"--measure-disk",
]
if known_at is not None:
argv += ["--known-at", str(known_at)]
Expand Down
29 changes: 21 additions & 8 deletions scripts/dev-fleet-warm-slot.py
Original file line number Diff line number Diff line change
Expand Up @@ -978,7 +978,7 @@ def warm(args: argparse.Namespace) -> dict[str, Any]:
tag = args.tag or f"warm-{args.slot}-{args.target[:8]}"
argv = build_command(checkout, tag, args.command)
env = build_env(derived, tag)
before_bytes = disk_bytes(layout.cache)
before_bytes = disk_bytes(layout.cache) if args.measure_disk else None
receipt = {
"schema_version": SCHEMA,
"kind": "warm",
Expand All @@ -995,9 +995,10 @@ def warm(args: argparse.Namespace) -> dict[str, Any]:
layout.logs / f"warm-{int(time.time())}-{log_token(args.target)}.log",
"warm", True, True, preempt_fd,
))
receipt["disk_bytes_before"] = before_bytes
receipt["disk_bytes_after"] = disk_bytes(layout.cache)
receipt["disk_growth_bytes"] = receipt["disk_bytes_after"] - before_bytes
if before_bytes is not None:
receipt["disk_bytes_before"] = before_bytes
receipt["disk_bytes_after"] = disk_bytes(layout.cache)
receipt["disk_growth_bytes"] = receipt["disk_bytes_after"] - before_bytes
atomic_json(layout.slot / "last-warm-receipt.json", receipt)
event(layout, "warm_finished", receipt=receipt)

Expand Down Expand Up @@ -1248,7 +1249,7 @@ def task_run(args: argparse.Namespace) -> dict[str, Any]:
match_class=match,
fallback_reason=fallback_reason,
):
before_bytes = disk_bytes(layout.cache)
before_bytes = disk_bytes(layout.cache) if args.measure_disk else None
build_started = time.time()
run = run_native(
layout, checkout, argv, env,
Expand Down Expand Up @@ -1276,11 +1277,13 @@ def task_run(args: argparse.Namespace) -> dict[str, Any]:
"toolchain": p["toolchain"],
"toolchain_fingerprint": p["toolchain_fingerprint"],
"derived_data_path": str(derived),
"disk_bytes_before": before_bytes,
}
if before_bytes is not None:
receipt["disk_bytes_before"] = before_bytes
receipt.update(run)
receipt["disk_bytes_after"] = disk_bytes(layout.cache)
receipt["disk_growth_bytes"] = receipt["disk_bytes_after"] - before_bytes
if before_bytes is not None:
receipt["disk_bytes_after"] = disk_bytes(layout.cache)
receipt["disk_growth_bytes"] = receipt["disk_bytes_after"] - before_bytes
receipt["source_after"] = head(checkout)
receipt["source_clean_after"] = clean(checkout)
try:
Expand Down Expand Up @@ -1457,6 +1460,11 @@ def make_parser() -> argparse.ArgumentParser:
p.add_argument("--owner", default="main-warmer")
p.add_argument("--tag")
p.add_argument("--ready-fd", type=int)
p.add_argument(
"--measure-disk",
action="store_true",
help="recursively measure cache bytes before/after native work; intended for benchmarks/diagnostics",
)
p.add_argument("command", nargs=argparse.REMAINDER)

p = sub.add_parser("task-base")
Expand All @@ -1481,6 +1489,11 @@ def make_parser() -> argparse.ArgumentParser:
p.add_argument("--tag")
p.add_argument("--known-at", type=float)
p.add_argument("--receipt", type=Path)
p.add_argument(
"--measure-disk",
action="store_true",
help="recursively measure cache bytes before/after native work; intended for benchmarks/diagnostics",
)
p.add_argument("command", nargs=argparse.REMAINDER)

p = sub.add_parser("release")
Expand Down
16 changes: 16 additions & 0 deletions tests/test_benchmark_dev_fleet_warm_slots.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import sys
import tempfile
import unittest
from unittest import mock

ROOT = Path(__file__).resolve().parents[1]
BENCH = ROOT / "scripts" / "benchmark-dev-fleet-warm-slots.py"
Expand Down Expand Up @@ -68,6 +69,21 @@ def test_wait_for_warmer_ready_uses_pipe_signal(self):
os.close(read_fd)
os.close(write_fd)

def test_benchmark_helpers_request_disk_measurement(self):
helper = self.root / "helper.py"
state = self.root / "state"
checkout = self.repo

with mock.patch.object(bench, "run_helper", return_value={"status": "ok"}) as run:
bench.warm(helper, state, checkout, "slot", self.main, [])
warm_argv = run.call_args.args[1]
self.assertIn("--measure-disk", warm_argv)

with mock.patch.object(bench, "run_helper", return_value={"status": "ok"}) as run:
bench.task(helper, state, checkout, "slot", self.main, "task", [])
task_argv = run.call_args.args[1]
self.assertIn("--measure-disk", task_argv)

def test_event_report_exposes_trial_metrics(self):
path = self.root / "events.jsonl"
rows = [
Expand Down
41 changes: 35 additions & 6 deletions tests/test_dev_fleet_warm_slot.py
Original file line number Diff line number Diff line change
Expand Up @@ -103,12 +103,12 @@ def call(self, *args, env=None, accepted=(0, 75, 130)):
def common(self, slot="slot"):
return ["--machine-state", str(self.state), "--slot", slot, "--checkout", str(self.repo)]

def warm(self, target, slot="slot", command=None, env=None):
return self.call(
"warm", *self.common(slot), "--target", target, "--",
*(command or native_command()),
env=env,
)
def warm(self, target, slot="slot", command=None, env=None, measure_disk=False):
argv = ["warm", *self.common(slot), "--target", target]
if measure_disk:
argv.append("--measure-disk")
argv += ["--", *(command or native_command())]
return self.call(*argv, env=env)

def task(
self,
Expand All @@ -118,8 +118,11 @@ def task(
command=None,
lease_id=None,
warm_generation_id=None,
measure_disk=False,
):
argv = ["task-run", *self.common(slot), "--target", target, "--task-id", task_id]
if measure_disk:
argv.append("--measure-disk")
if lease_id:
argv += ["--lease-id", lease_id]
if warm_generation_id:
Expand Down Expand Up @@ -217,6 +220,32 @@ def test_exact_generation_and_task_base(self):
planned = self.call("plan", *self.common(), "--target", self.base)
self.assertEqual(planned["reason"], "slot_needs_rewarm")

def test_recursive_disk_measurement_is_explicit(self):
default_warm = self.warm(self.base, slot="default-disk")
self.assertEqual(default_warm["status"], "warmed")
for field in ("disk_bytes_before", "disk_bytes_after", "disk_growth_bytes"):
self.assertNotIn(field, default_warm["receipt"])

measured_warm = self.warm(self.base, slot="measured-disk", measure_disk=True)
self.assertEqual(measured_warm["status"], "warmed")
for field in ("disk_bytes_before", "disk_bytes_after", "disk_growth_bytes"):
self.assertIn(field, measured_warm["receipt"])

default_task = self.task(self.base, slot="default-disk", task_id="default-task")
self.assertEqual(default_task["status"], "success")
for field in ("disk_bytes_before", "disk_bytes_after", "disk_growth_bytes"):
self.assertNotIn(field, default_task["receipt"])

measured_task = self.task(
self.base,
slot="measured-disk",
task_id="measured-task",
measure_disk=True,
)
self.assertEqual(measured_task["status"], "success")
for field in ("disk_bytes_before", "disk_bytes_after", "disk_growth_bytes"):
self.assertIn(field, measured_task["receipt"])

def test_reserved_generation_change_forces_cold_task_build(self):
warmed = self.warm(self.base)
selected = self.call(
Expand Down
Loading