diff --git a/client/campaign.py b/client/campaign.py index b0889360..0bb8ac29 100644 --- a/client/campaign.py +++ b/client/campaign.py @@ -275,16 +275,30 @@ def __init__(self, queue_dir: str, campaign_id: str, manifest: dict, max_storage raise ValueError("Campaign environment, recipes or protocol changed; start a new campaign") else: atomic_json(manifest_path, manifest) + # The retention allowance is campaign policy, not plan identity: it persists + # beside the manifest so resumes keep the original budget without making the + # frozen manifest drift for journals written before this separation existed. + budget_path = self.root / "budget.json" + # Entry points restore the saved allowance only when the caller has not + # supplied an explicit limit. Honor their resolved policy here, including + # deliberate increases or decreases, without changing campaign identity. + atomic_json(budget_path, {"schemaVersion": 1, "maxStorageMb": int(max_storage_mb)}) for path in sorted(self.root.glob("attempt-*.json")): record = load_record(json.loads(path.read_text())) info = record.metadata.get("info") or {} artifact = info.get("artifactPath") if artifact and not info.get("error"): candidate = Path(artifact) - if not candidate.is_file() or self.root.resolve() not in candidate.resolve().parents: + if candidate.is_file(): + if self.root.resolve() not in candidate.resolve().parents: + raise ValueError("Journal artifact missing or outside owned campaign") + if self.hash_file(candidate) != info.get("artifactSha256"): + raise ValueError("Journal artifact changed; cannot resume") + elif record.schedule.phase == "warmup": + if not self._warmup_released(record): + raise ValueError("Warmup artifact missing without verified release evidence") + elif not self.accepted_receipt(record): raise ValueError("Journal artifact missing or outside owned campaign") - if self.hash_file(candidate) != info.get("artifactSha256"): - raise ValueError("Journal artifact changed; cannot resume") self.records[record.schedule.execution_order] = record @staticmethod @@ -295,11 +309,97 @@ def hash_file(path): digest.update(chunk) return digest.hexdigest() + def _receipt_path(self, execution_order: int) -> Path: + return self.root / f"submission-{execution_order:06d}.accepted.json" + + def accepted_receipt(self, record): + """Return the accepted-upload receipt only when it is faithful to the attempt. + + A receipt counts when it names the same execution order, recipe, exact + artifact path and SHA-256 as the journaled record and carries the + server's run id. Corrupt, empty, partial or mismatched receipts never + authorize skipping an upload or trusting a released artifact.""" + info = record.metadata.get("info") or {} + artifact_path = str(info.get("artifactPath") or "").strip() + artifact_sha = info.get("artifactSha256") + if not artifact_path or not artifact_sha: + return None + try: + receipt = json.loads(self._receipt_path(record.schedule.execution_order).read_text()) + except (OSError, ValueError): + return None + if (isinstance(receipt, dict) + and receipt.get("schemaVersion") == 1 + and receipt.get("executionOrder") == record.schedule.execution_order + and receipt.get("recipeId") == record.schedule.recipe_id + and receipt.get("artifactPath") == artifact_path + and receipt.get("artifactSha256") == artifact_sha + and str(receipt.get("benchmarkRunId") or "").strip()): + return receipt + return None + + def _warmup_released(self, record) -> bool: + info = record.metadata.get("info") or {} + try: + evidence = json.loads((self.root / f"warmup-{record.schedule.execution_order:06d}.released.json").read_text()) + except (OSError, ValueError): + return False + return (isinstance(evidence, dict) + and evidence.get("schemaVersion") == 1 + and evidence.get("executionOrder") == record.schedule.execution_order + and evidence.get("recipeId") == record.schedule.recipe_id + and evidence.get("artifactPath") == str(info.get("artifactPath") or "").strip() + and evidence.get("artifactSha256") == info.get("artifactSha256")) + + def release_warmup_artifact(self, record) -> bool: + """Delete a hash-verified warmup output whose bytes have no consumer. + + Warmups tune encoders before measurement; they are never uploaded and + never re-read once the attempt (with its SHA-256) is durable. Release + evidence keeps the journal reopenable without the bytes; any hash + mismatch keeps the file in place so the budget check fails honestly.""" + info = record.metadata.get("info") or {} + artifact = str(info.get("artifactPath") or "").strip() + sha = info.get("artifactSha256") + if record.schedule.phase != "warmup" or not artifact or not sha: + return False + candidate = Path(artifact) + if self.root.resolve() not in candidate.resolve().parents or not candidate.is_file(): + return False + if self.hash_file(candidate) != sha: + return False + atomic_json(self.root / f"warmup-{record.schedule.execution_order:06d}.released.json", + {"schemaVersion": 1, + "executionOrder": record.schedule.execution_order, + "recipeId": record.schedule.recipe_id, + "artifactPath": artifact, + "artifactSha256": sha, + "releasedAt": time.time()}) + candidate.unlink() + return True + + def campaign_bytes(self): + return directory_bytes(str(self.root)) + def check_budget(self): - total = sum(path.stat().st_size for path in self.queue_root.rglob("*") if path.is_file()) - if total >= self.max_bytes: - raise OSError("Campaign storage budget reached; retained attempts can be resumed after freeing space") - return self.max_bytes - total + """This campaign's retention allowance plus a whole-volume safety floor. + + One campaign's retained bytes never block a different campaign: the + allowance below is charged to THIS journal only. The floor is global + protection for the disk itself and applies to every run.""" + used = self.campaign_bytes() + if used >= self.max_bytes: + raise OSError(f"this campaign retained {used // (1024 * 1024)} MB of its " + f"{self.max_bytes // (1024 * 1024)} MB allowance; accepted uploads retire " + f"automatically at checkpoints — wait for one, remove finished campaigns, " + f"or raise --max-storage-mb if the volume allows") + from shutil import disk_usage + floor_mb = max(0, int(os.environ.get("ENCODINGDB_MIN_FREE_MB", "1024"))) + free = disk_usage(str(self.queue_root)).free + if free < floor_mb * 1024 * 1024: + raise OSError(f"volume free space {free // (1024 * 1024)} MB is below the " + f"{floor_mb} MB safety floor reserved for the system; free space and retry") + return min(self.max_bytes - used, free - floor_mb * 1024 * 1024) def save(self, record): info = record.metadata.get("info") or {} @@ -364,3 +464,58 @@ def locked(): else: fcntl.flock(handle.fileno(), fcntl.LOCK_UN) return locked() + + +def active_collection(queue_dir: str): + """Describe the live collector holding this queue, if any. + + A held measurement lock (or a live encoder receipt during the brief + between-segment window) means an authoritative collection is running; a + second collector would corrupt measurement timing and is refused before + any expensive preparation.""" + root = Path(queue_dir) + lock = root / "measurement.lock" + if not lock.exists(): + return None + try: + with lock.open("a+b") as handle: + try: + if os.name == "nt": + import msvcrt + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) + msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) + else: + import fcntl + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + except OSError: + pass # Only a denied LOCK attempt means another collector holds it. + else: + return None + except OSError as exc: + # Opening/reading the lock is a permissions or media problem, NOT evidence + # of a running collector; report it accurately instead of mislabelling it. + raise OSError(f"Cannot inspect this queue's measurement lock: {exc}") from exc + import psutil + for receipt in sorted(root.glob("campaigns/*/*.active.json")): + try: + raw = json.loads(receipt.read_text()) + process = psutil.Process(int(raw.get("pid", -1))) + if process.create_time() == raw.get("createdAt"): + return {"campaignId": receipt.parent.name, "pid": int(raw["pid"])} + except Exception: + continue + return {"campaignId": None, "pid": None} + + +def directory_bytes(path: str) -> int: + """Bytes retained under a directory; a retiring upload is not an error.""" + total = 0 + for root, _dirs, names in os.walk(path): + for name in names: + try: + total += os.stat(os.path.join(root, name)).st_size + except FileNotFoundError: + continue + return total diff --git a/client/config.py b/client/config.py index 8939cd71..1cdc1548 100644 --- a/client/config.py +++ b/client/config.py @@ -70,6 +70,9 @@ def default_queue_dir() -> str: _BATCH_START_TS: float = 0.0 _BATCH_COMPLETED_COUNT: int = 0 +# New-attempt heartbeat across checkpoint segments; resumed attempts do not count. +_BATCH_ATTEMPTS_RECORDED: int = 0 + # Baseline cache for client-side outlier checks (populated lazily per session) _BASELINE_ROWS_CACHE: Optional[List[Dict[str, Any]]] = None _BASELINE_ROWS_CACHE_TS: float = 0.0 diff --git a/client/main.py b/client/main.py index 04c76d33..3d7bd1bc 100644 --- a/client/main.py +++ b/client/main.py @@ -11,6 +11,7 @@ import time import secrets from contextlib import nullcontext +from pathlib import Path from typing import Optional, Dict, Any, List, Tuple, Callable import psutil @@ -64,7 +65,7 @@ build_recipe_bootstrap, ) from .network import fetch_baseline_rows, check_compatibility -from .campaign import (CampaignJournal, atomic_json, physical_source_id, journal_path, +from .campaign import (CampaignJournal, active_collection, atomic_json, directory_bytes, physical_source_id, journal_path, PreparationScope, preparation_progress, check_preparation_cancelled, MeasurementBudget, MeasurementBudgetExceeded, check_measurement_budget, measurement_timeout, run_measurement_process) from .identity import selected_device @@ -77,6 +78,7 @@ RecipeSpec, StructuralExpectation, execute_protocol_campaign, + campaign_result_from_records, generate_campaign_id, ) from .spool import ( @@ -295,7 +297,12 @@ def progress(stage, **details): return label = details.get("clipId") or os.path.basename(str(details.get("path") or "")) done, total = details.get("completedBytes"), details.get("totalBytes") - amount = f" ({done}/{total} bytes)" if done is not None and total else "" + if done is not None and total: + amount = f" ({done}/{total} bytes)" + elif details.get("completed") is not None and details.get("total"): + amount = f" ({details['completed']}/{details['total']})" + else: + amount = "" print_info(f"Preparing: {stage} {label}{amount}") try: @@ -933,6 +940,20 @@ def _probe_artifact_contract(path: str) -> ArtifactProbe: ) +# Frozen suite sources are re-probed only when the file identity changes; the +# structural expectations for every recipe on one clip are byte-determined. +_SOURCE_CONTRACT_CACHE: Dict[str, ArtifactProbe] = {} + + +def _source_contract_cache_key(path: str) -> str: + resolved = os.path.realpath(path) + try: + stat = os.stat(resolved) + return f"{resolved}:{stat.st_size}:{stat.st_mtime_ns}" + except OSError: + return resolved + + def _build_protocol_config() -> ProtocolConfig: threshold = _safe_float(os.environ.get("ENCODINGDB_PROTOCOL_STABILITY_THRESHOLD")) adaptive_repeats = _safe_int(os.environ.get("ENCODINGDB_PROTOCOL_MAX_ADAPTIVE_REPEATS")) @@ -950,13 +971,33 @@ def _build_protocol_recipe_specs( default_input_hash: str, ) -> List[RecipeSpec]: specs: List[RecipeSpec] = [] - for task in tasks: + resolved_tasks = [ + (task, *_resolve_input_for_task(default_input_path, default_input_hash, task)) + for task in tasks + ] + contract_total = len({_source_contract_cache_key(item[1]) for item in resolved_tasks}) or 1 + contracts_done = 0 + for task, effective_input, input_hash, prepared_clip in resolved_tasks: + check_preparation_cancelled() encoder = str(task.get("encoder") or "").strip() preset = str(task.get("preset") or "").strip() crf = task.get("crf") rate_control = task.get("rateControl") - effective_input, input_hash, prepared_clip = _resolve_input_for_task(default_input_path, default_input_hash, task) - source_probe = _probe_artifact_contract(effective_input) + # One full-decode contract probe per distinct source file; every recipe + # on a frozen clip shares the same byte-determined expectations. + contract_key = _source_contract_cache_key(effective_input) + source_probe = _SOURCE_CONTRACT_CACHE.get(contract_key) + if source_probe is None: + preparation_progress( + "source-contract", + clipId=(prepared_clip.clip_id if isinstance(prepared_clip, PreparedSuiteClip) + else os.path.basename(effective_input)), + completed=contracts_done + 1, + total=contract_total, + ) + source_probe = _probe_artifact_contract(effective_input) + _SOURCE_CONTRACT_CACHE[contract_key] = source_probe + contracts_done += 1 source_duration = source_probe.duration_s source_fps = source_probe.avg_frame_rate source_frame_count = source_probe.frame_count @@ -1162,6 +1203,35 @@ def _json_default(value: Any) -> Any: return final_path +def _retire_uploaded_artifact(journal_root, record: Any, artifact_sha256: str, benchmark_run_id: str) -> None: + """Record the server-accepted receipt, then release the campaign copy. + + Accepted bytes now live on the server under ``benchmark_run_id``. The + journal re-opens behind this receipt only when it matches the attempt's + recipe, path and SHA-256, so evidence stays verifiable while the storage + budget only holds attempts that still need uploading.""" + if not str(benchmark_run_id or "").strip() or not artifact_sha256: + return # Without a server run id nothing may be treated as accepted. + info = record.metadata.get("info") or {} + artifact_path = str(info.get("artifactPath") or "").strip() + atomic_json(journal_root / f"submission-{record.schedule.execution_order:06d}.accepted.json", + {"schemaVersion": 1, + "executionOrder": record.schedule.execution_order, + "recipeId": record.schedule.recipe_id, + "artifactPath": artifact_path, + "artifactSha256": artifact_sha256, + "benchmarkRunId": str(benchmark_run_id).strip(), + "acceptedAt": time.time()}) + if artifact_path: + candidate = Path(artifact_path).resolve() + # Only bytes owned by this campaign journal may ever be released. + if Path(journal_root).resolve() in candidate.parents and candidate.is_file(): + try: + candidate.unlink() + except OSError: + pass + + def _submit_payload_with_spool( *, queue_dir: str, @@ -1248,6 +1318,7 @@ def _matching_saved_sweep_campaign(queue_dir: str, mode: str, if saved.get("sweepMode") != mode or saved.get("seed") is None: continue if saved.get("tasks") == planned_tasks: + saved["campaignId"] = campaign_id # informational; identity match stays task-exact return saved return None @@ -1266,6 +1337,9 @@ def run_sweep_mode( if mode not in sweep_plan.SWEEP_MODES: print(f"Unsupported sweep mode: {mode}", file=sys.stderr) return 4 + refused = active_collection_guard(base_args.queue_dir, event_sink, scope="sweep") + if refused: + return refused base_args = _apply_submission_policy(base_args, interactive=interactive) preflight_rc = _preparation_preflight(base_args, event_sink=event_sink) if preflight_rc: @@ -1312,7 +1386,55 @@ def run_sweep_mode( campaign_seed = saved_manifest["seed"] plan_metadata = {key: saved_manifest[key] for key in sweep_plan.MANIFEST_KEYS if key in saved_manifest} print_info(f"Continuing the retained {mode} sweep campaign from its saved plan.") + storage_mb = int(getattr(base_args, "max_storage_mb", 2048)) + storage_explicit = bool(getattr(base_args, "max_storage_mb_explicit", False)) + saved_storage = 0 + if saved_manifest is not None and not storage_explicit and saved_manifest.get("campaignId"): + try: + budget = json.loads((journal_path(base_args.queue_dir, saved_manifest["campaignId"]) + / "budget.json").read_text(encoding="utf-8")) + saved_storage = int(budget.get("maxStorageMb") or 0) + except (OSError, ValueError, json.JSONDecodeError): + saved_storage = 0 + if saved_storage > 0 and saved_storage != storage_mb: + storage_mb = saved_storage + print_info(f"Restoring this campaign's original {storage_mb} MB storage allowance for resume.") + retained = directory_bytes(str(journal_path(base_args.queue_dir, saved_manifest["campaignId"]))) \ + if saved_manifest and saved_manifest.get("campaignId") else 0 + if saved_manifest is not None and not storage_explicit and saved_storage == 0 and retained > 0: + # Journals created before budgets were persisted carry no budget.json; recover an + # allowance that at least fits what the campaign already retains plus headroom, + # instead of resuming into an instant budget rejection. + storage_mb = max(storage_mb, retained // (1024 * 1024) + 1024, 6144 if mode == "large" else 0) + print_info(f"No persisted budget for this campaign; sizing retention to its retained " + f"{retained // (1024 * 1024)} MB plus 1024 MB headroom; --max-storage-mb overrides.") + if not storage_explicit and saved_manifest is None and mode == "large" and storage_mb < 6144: + # One full measured round of the 357-group large plan retains ~2.5-3 GiB + # before the first checkpoint upload retires bytes (observed 2026-09-22); + # the 2048 MB default would stall mid-round. --max-storage-mb overrides. + storage_mb = 6144 + print_info("Large sweep retention default is 6144 MB (one full measured round needs ~2.5-3 GiB " + "before checkpoint uploads retire bytes); override with --max-storage-mb.") + from shutil import disk_usage as _disk_usage + floor_mb = max(0, int(os.environ.get("ENCODINGDB_MIN_FREE_MB", "1024"))) + try: + usable = _disk_usage(str(base_args.queue_dir)).free - floor_mb * 1024 * 1024 + retained + except OSError: + usable = None # Queue volume not readable yet; journal open enforces budgets at creation. protocol_config = _build_protocol_config() + # The allowance is a ceiling, not a reservation. Small sweeps need much + # less space than that ceiling; resumed campaigns can upload completed + # groups as they advance, so require working room beyond retained bytes. + estimate = len(tasks) * (protocol_config.minimum_measured_runs + protocol_config.max_adaptive_repeats) * 12 * 1024 * 1024 + required = retained + 64 * 1024 * 1024 if saved_manifest else min(storage_mb * 1024 * 1024, max(128 * 1024 * 1024, estimate)) + if usable is not None and usable < required: + message = (f"This volume cannot hold the sweep's estimated {required // (1024 * 1024)} MB working set: about " + f"{max(0, usable) // (1024 * 1024)} MB is usable beyond the {floor_mb} MB free-space " + f"floor. Free disk space, choose a smaller mode, or pass --max-storage-mb with a " + f"value the volume supports; nothing was started.") + print(message, file=sys.stderr) + _emit_event(event_sink, "run_error", scope="preparation", code=6, message=message) + return 6 per_recipe_min = protocol_config.warmup_runs + protocol_config.minimum_measured_runs per_recipe_max = per_recipe_min + protocol_config.max_adaptive_repeats encodes_min = len(tasks) * per_recipe_min @@ -1320,11 +1442,10 @@ def run_sweep_mode( if campaign_seed is None: env_seed = _safe_int(os.environ.get("ENCODINGDB_PROTOCOL_SEED")) campaign_seed = env_seed if env_seed is not None else secrets.randbits(63) - explicit_duration = bool(getattr(base_args, "explicit_max_duration_minutes", False)) + explicit_duration = bool(getattr(base_args, "max_duration_minutes_explicit", False)) segment_minutes = float(getattr(base_args, "max_duration_minutes", 60)) - attempts_cap = (int(getattr(base_args, "max_attempts")) if bool(getattr(base_args, "explicit_max_attempts", False)) + attempts_cap = (int(getattr(base_args, "max_attempts")) if bool(getattr(base_args, "max_attempts_explicit", False)) else encodes_max) - storage_mb = int(getattr(base_args, "max_storage_mb", 2048)) print_info( f"{mode} sweep: {len(plan.steps)} native recipes across {len(plan.encoders)} encoders " f"on {len(suite_clips)} frozen clip(s) = {len(tasks)} measured groups; " @@ -1348,6 +1469,7 @@ def run_sweep_mode( config._BATCH_ACTIVE = True config._BATCH_START_TS = time.perf_counter() config._BATCH_COMPLETED_COUNT = 0 + config._BATCH_ATTEMPTS_RECORDED = 0 total_submitted = 0 try: segment = 0 @@ -1382,10 +1504,11 @@ def run_sweep_mode( if rc == 11 and _is_cancelled(cancel_event): rc = 130 break - if int(getattr(config, "_BATCH_COMPLETED_COUNT", 0)) <= 0: + if int(getattr(config, "_BATCH_ATTEMPTS_RECORDED", 0)) <= 0: print("Checkpoint reached without any new measurement; stopping to keep retained progress " "safe. Start the same mode to continue.", file=sys.stderr) break + config._BATCH_ATTEMPTS_RECORDED = 0 if segment >= 10000: print("Safety segment limit reached; campaign remains saved and continues on the next start.", file=sys.stderr) @@ -1396,7 +1519,15 @@ def run_sweep_mode( elapsed_sec = max(0.0, time.perf_counter() - config._BATCH_START_TS) if show_end_screen: _clear_screen() - print_end_screen(total_submitted, elapsed_sec) + end_status = ("complete" if rc == 0 else + "paused" if rc in (10, 11) else + "interrupted" if rc == 130 else "failed") + print_end_screen(total_submitted, elapsed_sec, status=end_status, + recovery=None if rc == 0 else + ("Retained campaign saved; start this mode again to continue it." + if rc in (10, 11, 130) else + "Nothing was marked complete; fix the error above and start again — " + "the retained campaign continues from its journal.")) try: if os.name == "nt" and (bool(getattr(base_args, "pause_on_exit", False)) or bool(getattr(sys, "frozen", False))): input("Press Enter to exit...") @@ -1414,6 +1545,33 @@ def sweep_plan_label(encoder: str) -> str: return encoder +def active_collection_guard(queue_dir: str, + event_sink: Optional[Callable[[Dict[str, Any]], None]] = None, + *, scope: str = "batch") -> Optional[int]: + """Refusal code when this queue already hosts a live collection, else None. + + A live collector holds the measurement lock for its whole campaign: its + checkpoints continue automatically, so a second run must wait for it to + finish or stop/cancel it first - never "resume over" a checkpoint.""" + try: + active = active_collection(str(queue_dir)) + except OSError as exc: + message = str(exc) + print(message, file=sys.stderr) + _emit_event(event_sink, "run_error", scope=scope, code=6, message=message) + return 6 + if active is None: + return None + who = f" (campaign {active['campaignId']}, PID {active['pid']})" if active.get("campaignId") else "" + message = (f"Another collection is actively running in this queue{who}. Its checkpoints continue " + f"automatically - let it finish, or stop/cancel that run first; a second collector " + f"would corrupt measurement timing.") + print(message, file=sys.stderr) + _emit_event(event_sink, "run_error", scope=scope, code=6, message=message) + return 6 + + + @_preparation_operation def run_benchmark_batch( *, @@ -1425,12 +1583,16 @@ def run_benchmark_batch( cancel_event: Optional[Any] = None, plan_metadata: Optional[Dict[str, Any]] = None, ) -> int: + refused = active_collection_guard(args.queue_dir, event_sink, scope="batch") + if refused: + return refused duration_minutes = float(getattr(args, "max_duration_minutes", 60)) if not math.isfinite(duration_minutes) or not math.isfinite(duration_minutes * 60) or duration_minutes <= 0: message = "--max-duration-minutes must be positive and finite" print(message, file=sys.stderr) _emit_event(event_sink, "run_error", scope="batch", code=4, message=message) return 4 + campaign_paused = False preflight_rc = _preparation_preflight(args, base_url=base_url, event_sink=event_sink) if preflight_rc: return preflight_rc @@ -1494,12 +1656,25 @@ def run_benchmark_batch( manifest.update(plan_metadata) try: journal = CampaignJournal(args.queue_dir, campaign_id, manifest, int(getattr(args, "max_storage_mb", 2048))) - journal.check_budget() + remaining_bytes = journal.check_budget() except Exception as exc: - print(f"Cannot open campaign journal: {exc}", file=sys.stderr) + message = f"Cannot open campaign journal: {exc}" + print(message, file=sys.stderr) _debug_exception_traceback() + _emit_event(event_sink, "run_error", scope="batch", code=6, message=message) return 6 print_info(f"Campaign {campaign_id}: at most {total_tasks} encodes; resume with --resume-campaign {campaign_id}") + pending_attempts = max(0, total_tasks - len(journal.records)) + if pending_attempts: + observed = sorted(int((record.metadata.get("info") or {}).get("fileSizeBytes") or 0) + for record in journal.records.values()) + typical = observed[len(observed) // 2] if any(observed) else 12 * 1024 * 1024 + needed = pending_attempts * typical + if needed > remaining_bytes: + print_info(f"Planned attempts could need up to ≈{needed // (1024 * 1024)} MB of retention while " + f"this campaign's remaining allowance is {remaining_bytes // (1024 * 1024)} MB. " + f"Completed groups upload and retire automatically at checkpoints; if the volume " + f"allows, --max-storage-mb raises the allowance.") total_batches = 1 run_started_at = time.perf_counter() use_token = _should_use_submit_token(args) @@ -1756,15 +1931,37 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: ) budget = MeasurementBudget(duration_minutes, cancel_event=cancel_event) - with budget.activate(): - campaign_result = execute_protocol_campaign( - recipes=recipe_specs, + recorded_before = set(journal.records) + + def _journal_attempt(record: Any) -> None: + journal.save(record) + if record.schedule.execution_order not in recorded_before: + with config._GLOBAL_STATE_LOCK: + config._BATCH_ATTEMPTS_RECORDED += 1 + journal.release_warmup_artifact(record) # Hash-verified, never publishable. + + try: + with budget.activate(): + campaign_result = execute_protocol_campaign( + recipes=recipe_specs, + config=protocol_config, + encode_runner=_encode_protocol_run, + environment_sampler=_sample_environment, + seed=campaign_seed, + record_sink=_journal_attempt, + resumed_records=journal.records, + ) + except MeasurementBudgetExceeded: + # The 60-minute checkpoint is a resumable pause: every attempt + # recorded so far is journalized, so rebuild the scheduler's + # view and upload only terminal measurement groups below. + campaign_paused = True + campaign_result = campaign_result_from_records( + campaign_id=campaign_id, config=protocol_config, - encode_runner=_encode_protocol_run, - environment_sampler=_sample_environment, seed=campaign_seed, - record_sink=journal.save, - resumed_records=journal.records, + recipes=recipe_specs, + records=journal.records, ) attempt_evidence_path = _persist_protocol_attempt_evidence(args.queue_dir, campaign_result) _emit_event( @@ -1777,8 +1974,11 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: measurement_groups = _completed_measurement_groups(campaign_result) measured_records: List[Tuple[RecipeSpec, Any]] = [] + unfinished_recipes = getattr(campaign_result, "unfinished_recipes", frozenset()) for recipe_result in campaign_result.recipe_results: recipe = recipe_by_id[recipe_result.recipe_id] + if recipe_result.recipe_id in unfinished_recipes: + continue # A later segment can still extend a checkpointed group. _emit_event( event_sink, "protocol_recipe_complete", @@ -1796,6 +1996,8 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: if not getattr(args, "local_metrics", False): record.metadata["metrics"] = {} continue + if journal.accepted_receipt(record) is not None: + continue # Bytes retired after acceptance; nothing left to measure. if _is_cancelled(cancel_event): raise KeyboardInterrupt info = dict(record.metadata.get("info") or {}) @@ -1855,6 +2057,8 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: for recipe, record in measured_records: if _is_cancelled(cancel_event): raise KeyboardInterrupt + if journal.accepted_receipt(record) is not None: + continue # Faithful accepted receipt from an earlier segment. task = _task_from_recipe(recipe) info = dict(record.metadata.get("info") or {}) codec_label = str(info.get('encoderUsed') or task['encoder']) @@ -2125,6 +2329,7 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: ) if status == "submitted": submitted_count += 1 + _retire_uploaded_artifact(journal.root, record, artifact_sha256, error_text) if error_text: print_info(f"Authoritative benchmark run recorded as {error_text}.") _emit_event( @@ -2209,6 +2414,23 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: processed_total += 1 progress.advance(description=_batch_status("Completed", processed_total, str(payload['codec']), str(payload['preset']))) _emit_event(event_sink, "task_complete", scope="batch", processed=processed_total, total=total_tasks) + if campaign_paused: + status = {"status": "budget_exhausted", "campaignId": campaign_id, + "maxDurationMinutes": duration_minutes, + "elapsedSeconds": max(0.0, budget.clock() - budget.started), + "stoppedAt": time.time(), "phase": "measurement", + "checkpointUploads": submitted_count, "queuedUploads": queued_count} + flight = journal.root / "in-flight.json" + if flight.exists(): + try: + status["lastStartedAttempt"] = json.loads(flight.read_text()) + except (OSError, ValueError): + pass + atomic_json(journal.root / f"budget-exhausted-{time.time_ns()}.json", status) + print_warning(f"Checkpoint reached: {submitted_count} accepted upload(s) this segment; the " + f"retained campaign continues automatically ({campaign_id}).") + _emit_event(event_sink, "run_budget_exhausted", scope="batch", **status) + return 11 except MeasurementBudgetExceeded as exc: status = {"status": "budget_exhausted", "campaignId": campaign_id, "maxDurationMinutes": exc.budget.minutes, @@ -2226,6 +2448,16 @@ def _encode_protocol_run(schedule: Any, recipe: RecipeSpec) -> EncodeOutcome: print_warning(f"{exc}. Saved attempts remain available; resume with --resume-campaign {campaign_id}.") _emit_event(event_sink, "run_budget_exhausted", scope="batch", **status) return 11 + except SpoolCapacityError as exc: + if campaign_paused: + print_warning(f"Checkpoint upload deferred: {exc}; the retained campaign continues.") + _emit_event(event_sink, "run_budget_exhausted", scope="batch", + status="budget_exhausted", campaignId=campaign_id, deferred=str(exc)) + return 11 + print(f"Campaign retained for resume: {exc}", file=sys.stderr) + _debug_exception_traceback() + _emit_event(event_sink, "run_error", scope="batch", code=6, message=str(exc)) + return 6 except (OSError, ValueError, TimeoutError) as exc: print(f"Campaign retained for resume: {exc}", file=sys.stderr) _debug_exception_traceback() @@ -2372,6 +2604,9 @@ def run_v7_suite_clip_mode( @_preparation_operation def _resume_campaign(args, *, event_sink=None, cancel_event=None, interactive=False): + refused = active_collection_guard(args.queue_dir, event_sink, scope="resume") + if refused: + return refused args = _apply_submission_policy(args, interactive=interactive) preflight_rc = _preparation_preflight(args, event_sink=event_sink) if preflight_rc: @@ -2380,16 +2615,55 @@ def _resume_campaign(args, *, event_sink=None, cancel_event=None, interactive=Fa root = journal_path(args.queue_dir, args.resume_campaign) saved = json.loads((root / "manifest.json").read_text()) args.campaign_seed = saved["seed"] + try: + budget = json.loads((root / "budget.json").read_text(encoding="utf-8")) + persisted_mb = int(budget.get("maxStorageMb") or 0) + except (OSError, ValueError, json.JSONDecodeError): + persisted_mb = 0 + if not bool(getattr(args, "max_storage_mb_explicit", False)): + if persisted_mb > 0: + args.max_storage_mb = persisted_mb + print_info(f"Restoring this campaign's original {persisted_mb} MB storage allowance for resume.") + else: + retained_mb = directory_bytes(str(root)) // (1024 * 1024) + args.max_storage_mb = max(int(getattr(args, "max_storage_mb", 2048)), retained_mb + 1024, + 6144 if saved.get("sweepMode") == "large" else 0) + print_info(f"No saved storage allowance; allowing {args.max_storage_mb} MB for this " + "campaign's retained files and continuation.") # Reopening the journal requires the exact saved manifest; sweep campaigns persist # their planner metadata, so resume must pass every plan key through unchanged. plan_metadata = {key: saved[key] for key in sweep_plan.MANIFEST_KEYS if key in saved} or None - tasks = [{"encoder": task["encoder"], "preset": task["preset"], "crf": task["crf"], - "rateControl": task["rateControl"], "suiteClip": _prepare_named_suite_clip(task["clipId"])} - for task in saved["tasks"]] + clips = {} + tasks = [] + for task in saved["tasks"]: + check_preparation_cancelled() + clip_id = task["clipId"] + if clip_id not in clips: + clips[clip_id] = _prepare_named_suite_clip(clip_id) + tasks.append({"encoder": task["encoder"], "preset": task["preset"], "crf": task["crf"], + "rateControl": task["rateControl"], "suiteClip": clips[clip_id]}) check_preparation_cancelled() - return run_benchmark_batch(hardware=detect_hardware(), base_url=args.base_url, args=args, tasks=tasks, - event_sink=event_sink, cancel_event=cancel_event, - plan_metadata=plan_metadata) + protocol_config = _build_protocol_config() + if not bool(getattr(args, "max_attempts_explicit", False)): + args.max_attempts = max(int(getattr(args, "max_attempts", 100)), len(tasks) * ( + protocol_config.warmup_runs + protocol_config.minimum_measured_runs + protocol_config.max_adaptive_repeats)) + hardware = detect_hardware() + continuing_sweep = saved.get("sweepMode") in sweep_plan.SWEEP_MODES + for _segment in range(10000): + config._BATCH_ATTEMPTS_RECORDED = 0 + rc = run_benchmark_batch(hardware=hardware, base_url=args.base_url, args=args, tasks=tasks, + event_sink=event_sink, cancel_event=cancel_event, + plan_metadata=plan_metadata) + if _is_cancelled(cancel_event): + return 130 + if rc != 11 or not continuing_sweep or bool(getattr(args, "max_duration_minutes_explicit", False)): + return rc + if config._BATCH_ATTEMPTS_RECORDED <= 0: + print_warning("Checkpoint made no new measurement; retained campaign is paused safely.") + return 11 + print_info("Continuing the saved sweep after its checkpoint.") + print_warning("Safety segment limit reached; the campaign remains saved for continuation.") + return 11 except Exception as exc: print(f"Cannot resume campaign: {exc}", file=sys.stderr) _debug_exception_traceback() @@ -2976,9 +3250,13 @@ def build_arg_parser() -> argparse.ArgumentParser: p.add_argument("--max-duration-minutes", type=float, default=60, action=_ExplicitBudgetAction, help="Measurement allowance per invocation in minutes; acquisition and uploads are " "separate (default 60; guided sweeps continue across checkpoints unless set)") - p.add_argument("--max-storage-mb", type=int, default=2048, help="Maximum retained queue and campaign storage in MiB") p.add_argument("--legacy-diagnostic", action="store_true", help="Noncanonical local-only legacy diagnostic; never publishes") - p.set_defaults(explicit_max_attempts=False, explicit_max_duration_minutes=False) + p.add_argument("--max-storage-mb", type=int, default=2048, action=_ExplicitBudgetAction, + help="Maximum retained queue and campaign storage in MiB (default 2048; the large " + "sweep defaults to 6144 because one measured round retains ~2.5-3 GiB before " + "checkpoint uploads retire bytes)") + p.set_defaults(max_attempts_explicit=False, max_duration_minutes_explicit=False, + max_storage_mb_explicit=False) return p diff --git a/client/protocol.py b/client/protocol.py index e903a8ca..a40eaf22 100644 --- a/client/protocol.py +++ b/client/protocol.py @@ -390,6 +390,12 @@ class CampaignResult: protocol_version: str seed: Optional[int] recipe_results: List[RecipeCampaignResult] + # A budget-exhausted segment is a resumable pause, not a failed campaign: + # completed=False carries the pause through submission, and unfinished + # recipes (still short of a terminal stability outcome) stay unsubmitted + # so their measurement groups can complete in a later segment untouched. + completed: bool = True + unfinished_recipes: frozenset = frozenset() def to_dict(self) -> Dict[str, Any]: return { @@ -397,6 +403,8 @@ def to_dict(self) -> Dict[str, Any]: "protocolVersion": self.protocol_version, "seed": self.seed, "recipeResults": [result.to_dict() for result in self.recipe_results], + "completed": self.completed, + "unfinishedRecipes": sorted(self.unfinished_recipes), } @@ -1113,3 +1121,50 @@ def environment_sampler(schedule, recipe): seed=effective_seed, recipe_results=recipe_results, ) + + +def campaign_result_from_records( + *, + campaign_id: str, + config: ProtocolConfig, + seed: Optional[int], + recipes: Sequence[RecipeSpec], + records: Dict[int, BenchmarkRunRecord], +) -> CampaignResult: + """Rebuild the campaign view from journalized attempts at a checkpoint. + + The terminal rule mirrors the scheduler's own exits: a recipe is final + exactly when its measured group is stable or has reached the measured + attempt cap. Everything else stays in ``unfinished_recipes`` so a + checkpoint never uploads a group a later segment could still extend.""" + by_recipe: Dict[str, List[BenchmarkRunRecord]] = {recipe.recipe_id: [] for recipe in recipes} + for record in sorted(records.values(), key=lambda item: item.schedule.execution_order): + if record.schedule.campaign_id == campaign_id and record.schedule.recipe_id in by_recipe: + by_recipe[record.schedule.recipe_id].append(record) + attempt_cap = config.minimum_measured_runs + config.max_adaptive_repeats + recipe_results: List[RecipeCampaignResult] = [] + unfinished = set() + for recipe in recipes: + runs = by_recipe[recipe.recipe_id] + stability = evaluate_stability(runs, config) + measured_runs = [run for run in runs if run.schedule.phase == "measured"] + if not stability.stable and len(measured_runs) < attempt_cap: + unfinished.add(recipe.recipe_id) + recipe_results.append( + RecipeCampaignResult( + recipe_id=recipe.recipe_id, + runs=runs, + stability=stability, + measured_runs_required=config.minimum_measured_runs, + measured_runs_completed=len(measured_runs), + measured_runs_counted=stability.sample_count, + ) + ) + return CampaignResult( + campaign_id=campaign_id, + protocol_version=config.version, + seed=seed, + recipe_results=recipe_results, + completed=False, + unfinished_recipes=frozenset(unfinished), + ) diff --git a/client/spool.py b/client/spool.py index 13cbb0c0..0e551dfb 100644 --- a/client/spool.py +++ b/client/spool.py @@ -11,6 +11,7 @@ from typing import Any, Dict, List, Optional, Tuple from .artifacts import AUTHORITATIVE_ARTIFACT_SUBMISSION_KIND, submit_artifact_submission +from .campaign import directory_bytes from .network import SubmitError, submit SPOOL_VERSION = 1 @@ -50,9 +51,25 @@ def _spool_write_lock(queue_dir: str): release() +def _publication_bytes(queue_dir: str) -> int: + """Shared publication store: pending entries, managed copies, receipts, dead + letters. OTHER campaigns' retained attempts are charged to their own journal + allowance, never to this run's budget.""" + def fail_scan(error): + raise error + total = 0 + for root, dirs, names in os.walk(queue_dir, onerror=fail_scan): + if os.path.abspath(root) == os.path.abspath(queue_dir) and "campaigns" in dirs: + dirs.remove("campaigns") + # Stat errors fail closed rather than undercount. + total += sum(os.stat(os.path.join(root, name)).st_size for name in names) + return total + + def _check_spool_capacity(queue_dir: str, payload: Dict[str, Any], max_storage_mb: int) -> None: staged = dict(payload) copy_bytes = 0 + source = "" if payload.get("submissionKind") == AUTHORITATIVE_ARTIFACT_SUBMISSION_KIND: source = str(payload.get("artifactPath") or "").strip() sha = str(payload.get("artifactSha256") or "").strip().lower() @@ -65,17 +82,31 @@ def _check_spool_capacity(queue_dir: str, payload: Dict[str, Any], max_storage_m copy_bytes = size staged.update(artifactPath=destination, artifactManaged=True) metadata_bytes = len(json.dumps(_envelope_for_payload(staged), sort_keys=True).encode("utf-8")) - # Count campaign originals, managed copies, pending records, receipts, terminal - # evidence and temporary files. Stat errors fail closed rather than undercount. - def fail_scan(error): - raise error - used = sum(os.stat(os.path.join(root, name)).st_size - for root, _, names in os.walk(queue_dir, onerror=fail_scan) for name in names) + used = _publication_bytes(queue_dir) + # Submission envelopes do not carry a top-level campaignId. Account for + # the original using its owned filesystem location, including older flat + # campaign layouts, instead of trusting optional payload metadata. + if source: + campaigns = Path(queue_dir).resolve() / "campaigns" + try: + relative = Path(source).resolve().relative_to(campaigns) + except (OSError, ValueError): + relative = None + if relative is not None and relative.parts: + owned = campaigns / relative.parts[0] + used += directory_bytes(str(owned)) if owned.is_dir() else owned.stat().st_size required = copy_bytes + metadata_bytes + SPOOL_METADATA_RESERVE_BYTES if used + required > max_storage_mb * 1024 * 1024: - raise SpoolCapacityError("Publication storage budget reached; increase --max-storage-mb and resume the retained campaign") - if shutil.disk_usage(queue_dir).free < required: - raise SpoolCapacityError("Insufficient free disk for publication staging; free space and resume the retained campaign") + raise SpoolCapacityError( + f"Publication storage budget reached: this campaign's retained attempts plus pending " + f"uploads exceed its {max_storage_mb} MB allowance ({used // (1024 * 1024)} MB used). " + f"Accepted uploads retire at checkpoints — wait for one, or raise --max-storage-mb if the volume allows") + floor_mb = max(0, int(os.environ.get("ENCODINGDB_MIN_FREE_MB", "1024"))) + free = shutil.disk_usage(queue_dir).free + if free - required < floor_mb * 1024 * 1024: + raise SpoolCapacityError(f"Insufficient free disk for publication staging: only " + f"{free // (1024 * 1024)} MB free against the {floor_mb} MB safety " + f"floor; free space and resume the retained campaign") @dataclass diff --git a/client/tests/test_campaign_durability.py b/client/tests/test_campaign_durability.py index 4173e502..b8c45609 100644 --- a/client/tests/test_campaign_durability.py +++ b/client/tests/test_campaign_durability.py @@ -273,3 +273,23 @@ def launch(*args, **kwargs): if process.poll() is None: process.kill() process.wait() + + +def test_per_campaign_allowance_isolates_other_retained_campaigns(tmp_path, monkeypatch): + monkeypatch.setenv("ENCODINGDB_MIN_FREE_MB", "1") + a = CampaignJournal(str(tmp_path), "campaign-" + "a" * 16, {"physicalSourceId": "TEST ONLY"}, 2048) + with open(a.root / "bulk.bin", "wb") as bulk: # Sparse stand-in for retained artifact bytes + bulk.truncate(2100 * 1024 * 1024) + with pytest.raises(OSError, match="this campaign retained"): + a.check_budget() + b = CampaignJournal(str(tmp_path), "campaign-" + "b" * 16, {"physicalSourceId": "TEST ONLY"}, 2048) + assert b.check_budget() > 0 # Another campaign's retention must not block a fresh plan + + +def test_disk_guard_reports_volume_free_space_floor(tmp_path, monkeypatch): + import shutil + journal = CampaignJournal(str(tmp_path), "campaign-" + "c" * 16, {"physicalSourceId": "TEST ONLY"}, 2048) + monkeypatch.setattr(shutil, "disk_usage", + lambda path: SimpleNamespace(total=1024 ** 4, free=500 * 1024 * 1024, used=0)) + with pytest.raises(OSError, match="safety floor"): + journal.check_budget() diff --git a/client/tests/test_main_routing.py b/client/tests/test_main_routing.py index f2668cab..42d63c44 100644 --- a/client/tests/test_main_routing.py +++ b/client/tests/test_main_routing.py @@ -948,9 +948,9 @@ def fake_batch(**kwargs): return 11 with tempfile.TemporaryDirectory() as queue_dir: - captured_args = self.args(queue_dir=queue_dir, max_duration_minutes=15.0, - explicit_max_duration_minutes=True, - max_attempts=999, explicit_max_attempts=True) + captured_args = self.args(queue_dir=queue_dir, + max_duration_minutes=15.0, max_duration_minutes_explicit=True, + max_attempts=999, max_attempts_explicit=True) with ExitStack() as stack: for patcher in self._sweep_run_patches(queue_dir, {}): stack.enter_context(patcher) @@ -964,11 +964,11 @@ def fake_batch(**kwargs): def test_sweep_auto_continues_across_checkpoints(self) -> None: calls = [] - client_main.config._BATCH_COMPLETED_COUNT = 0 + client_main.config._BATCH_ATTEMPTS_RECORDED = 0 def fake_batch(**kwargs): calls.append(kwargs) - client_main.config._BATCH_COMPLETED_COUNT += 1 + client_main.config._BATCH_ATTEMPTS_RECORDED += 1 return 11 if len(calls) == 1 else 0 try: @@ -980,7 +980,7 @@ def fake_batch(**kwargs): rc = client_main.run_sweep_mode(mode="small", base_args=self.args(queue_dir=queue_dir), show_end_screen=False, interactive=False, presets_cfg={}) finally: - client_main.config._BATCH_COMPLETED_COUNT = 0 + client_main.config._BATCH_ATTEMPTS_RECORDED = 0 self.assertEqual(rc, 0, "checkpointing continues until the plan completes") self.assertEqual(len(calls), 2) self.assertEqual(calls[0]["args"].campaign_seed, calls[1]["args"].campaign_seed) @@ -1030,6 +1030,185 @@ def fake_batch(**kwargs): self.assertEqual(rc, 130) self.assertEqual(called, [], "cancelled preparation must not start encoding") + def test_active_collection_refuses_second_collector(self) -> None: + import fcntl + with tempfile.TemporaryDirectory() as queue_dir: + with open(os.path.join(queue_dir, "measurement.lock"), "a+b") as held: + fcntl.flock(held.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + self.assertIsNotNone(client_main.active_collection(queue_dir)) + events = [] + rc = client_main.run_benchmark_batch( + hardware=HardwareInfo("CPU", "GPU", 16, "TestOS"), base_url="https://example.invalid", + args=self.args(queue_dir=queue_dir), tasks=[], + event_sink=lambda event: events.append(event)) + self.assertEqual(rc, 6) + self.assertTrue(any(event.get("type") == "run_error" and "actively running" in str(event.get("message")) + for event in events), "the GUI needs the honest refusal reason, not just exit 6") + + def test_failed_sweep_end_screen_reports_failure(self) -> None: + with tempfile.TemporaryDirectory() as queue_dir: + with ExitStack() as stack: + for patcher in self._sweep_run_patches(queue_dir, {}): + stack.enter_context(patcher) + stack.enter_context(mock.patch.object(client_main, "run_benchmark_batch", return_value=6)) + screens = [] + stack.enter_context(mock.patch.object(client_main, "print_end_screen", + side_effect=lambda *a, **k: screens.append(k))) + rc = client_main.run_sweep_mode(mode="small", base_args=self.args(queue_dir=queue_dir), + show_end_screen=True, interactive=False, presets_cfg={}) + self.assertEqual(rc, 6) + self.assertEqual(screens[0].get("status"), "failed") + self.assertTrue(screens[0].get("recovery"), "a failed run must tell the operator how to continue") + + def test_sweep_entry_refuses_before_any_expensive_work(self) -> None: + import fcntl + with tempfile.TemporaryDirectory() as queue_dir: + with open(os.path.join(queue_dir, "measurement.lock"), "a+b") as held: + fcntl.flock(held.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + events = [] + rc = client_main.run_sweep_mode( + mode="small", base_args=self.args(queue_dir=queue_dir), + event_sink=lambda event: events.append(event), + show_end_screen=False, interactive=False, presets_cfg={}) + self.assertEqual(rc, 6) + self.assertTrue(any(event.get("type") == "run_error" and "actively running" in str(event.get("message")) + and "continue automatically" in str(event.get("message")) + for event in events), + "refusal must explain auto-continuing checkpoints, not suggest resuming over them") + + @unittest.skipIf(os.name != "posix" or hasattr(os, "geteuid") and os.geteuid() == 0, + "permission semantics require a non-root POSIX account") + def test_active_probe_reports_lock_permission_error_accurately(self) -> None: + with tempfile.TemporaryDirectory() as queue_dir: + lock = os.path.join(queue_dir, "measurement.lock") + open(lock, "a+b").close() + os.chmod(lock, 0) + try: + with self.assertRaisesRegex(OSError, "Cannot inspect"): + client_main.active_collection(queue_dir) + finally: + os.chmod(lock, 0o644) + + def test_sweep_restores_saved_storage_allowance_on_resume(self) -> None: + clip = MainRoutingTests("_quick_clip")._quick_clip() + plan = client_main.sweep_plan.plan_sweep("small", ["libx264"], presets_cfg={}) + saved_tasks = [{"encoder": step["encoder"], "preset": step["preset"], "crf": step["crf"], + "rateControl": step["rateControl"], "clipId": clip.clip_id} for step in plan.steps] + with tempfile.TemporaryDirectory() as queue_dir: + self._saved_sweep_journal(queue_dir, saved_tasks) + journal_root = client_main.journal_path(queue_dir, "campaign-0123456789abcdef") + (journal_root / "budget.json").write_text(json.dumps({"schemaVersion": 1, "maxStorageMb": 6144})) + captured = {} + from types import SimpleNamespace + with ExitStack() as stack: + for patcher in self._sweep_run_patches(queue_dir, captured): + stack.enter_context(patcher) + stack.enter_context(mock.patch("shutil.disk_usage", + return_value=SimpleNamespace(total=1024 ** 4, used=0, + free=100 * 1024 ** 3))) + rc = client_main.run_sweep_mode(mode="small", base_args=self.args(queue_dir=queue_dir), + show_end_screen=False, interactive=False, presets_cfg={}) + self.assertEqual(rc, 0) + self.assertEqual(captured["args"].max_storage_mb, 6144, + "a campaign started with 6144 MB must not resume at the 2048 default and reject its own 2.7 GiB") + + def test_legacy_journal_without_budget_recovers_allowance(self) -> None: + from types import SimpleNamespace + clip = MainRoutingTests("_quick_clip")._quick_clip() + plan = client_main.sweep_plan.plan_sweep("small", ["libx264"], presets_cfg={}) + saved_tasks = [{"encoder": step["encoder"], "preset": step["preset"], "crf": step["crf"], + "rateControl": step["rateControl"], "clipId": clip.clip_id} for step in plan.steps] + with tempfile.TemporaryDirectory() as queue_dir: + self._saved_sweep_journal(queue_dir, saved_tasks) + journal_root = client_main.journal_path(queue_dir, "campaign-0123456789abcdef") + with open(journal_root / "retained.bin", "wb") as bulk: + bulk.truncate(2100 * 1024 * 1024) + captured = {} + with ExitStack() as stack: + for patcher in self._sweep_run_patches(queue_dir, captured): + stack.enter_context(patcher) + stack.enter_context(mock.patch("shutil.disk_usage", + return_value=SimpleNamespace(total=1024 ** 4, used=0, + free=100 * 1024 ** 3))) + rc = client_main.run_sweep_mode(mode="small", base_args=self.args(queue_dir=queue_dir), + show_end_screen=False, interactive=False, presets_cfg={}) + self.assertEqual(rc, 0) + self.assertEqual(captured["args"].max_storage_mb, 3124, + "a pre-budget.json campaign retaining 2100 MB must resume above its own retained bytes") + + def test_large_mode_default_storage_and_explicit_override(self) -> None: + plan = client_main.sweep_plan.plan_sweep("small", ["libx264"], presets_cfg={}) + from types import SimpleNamespace + for changes, expected in (({}, 6144), ({"max_storage_mb": 512, "max_storage_mb_explicit": True}, 512)): + with self.subTest(**changes), mock.patch("shutil.disk_usage") as usage: + usage.return_value = SimpleNamespace(total=1024 ** 4, used=0, free=100 * 1024 ** 3) + with tempfile.TemporaryDirectory() as queue_dir: + captured = {} + with ExitStack() as stack: + for patcher in self._sweep_run_patches(queue_dir, captured): + stack.enter_context(patcher) + stack.enter_context(mock.patch.object(client_main.sweep_plan, "plan_sweep", + return_value=plan)) + rc = client_main.run_sweep_mode( + mode="large", base_args=self.args(queue_dir=queue_dir, **changes), + show_end_screen=False, interactive=False, presets_cfg={}) + self.assertEqual(rc, 0) + self.assertEqual(captured["args"].max_storage_mb, expected) + + def test_volume_cannot_host_working_set_stops_before_any_encode(self) -> None: + from types import SimpleNamespace + calls = [] + with tempfile.TemporaryDirectory() as queue_dir: + with ExitStack() as stack: + for patcher in self._sweep_run_patches(queue_dir, {}): + stack.enter_context(patcher) + stack.enter_context(mock.patch("shutil.disk_usage", + return_value=SimpleNamespace(total=1024 ** 4, used=0, free=1088 * 1024 * 1024))) + events = [] + rc = client_main.run_sweep_mode( + mode="small", base_args=self.args(queue_dir=queue_dir, max_storage_mb=6144, + max_storage_mb_explicit=True), + event_sink=lambda event: events.append(event), + show_end_screen=False, interactive=False, presets_cfg={}) + self.assertEqual(rc, 6) + self.assertEqual(calls, [], "a plan the volume cannot host must never reach encoding") + self.assertTrue(any(event.get("type") == "run_error" and "cannot hold" in str(event.get("message")) + for event in events)) + + def test_end_screen_rendered_text_is_honest(self) -> None: + import contextlib + import io + from client import ui + for status, needle in (("complete", "Benchmark complete"), + ("paused", "not finished"), + ("interrupted", "not finished"), + ("failed", "nothing was marked finished")): + buf = io.StringIO() + with contextlib.redirect_stdout(buf): + ui.print_end_screen(7, 83, status=status, recovery="Recovery hint.") + text = buf.getvalue() + self.assertIn(needle, text, f"{status} end screen") + self.assertIn("Recovery hint.", text) + buf = io.StringIO() + with contextlib.redirect_stdout(buf): + ui.print_end_screen(3, 5, status="failed", recovery="x") + self.assertNotIn("Benchmark complete", buf.getvalue()) + + def test_failed_sweep_prints_failed_banner_without_mock(self) -> None: + import contextlib + import io + with tempfile.TemporaryDirectory() as queue_dir: + with ExitStack() as stack: + for patcher in self._sweep_run_patches(queue_dir, {}): + stack.enter_context(patcher) + stack.enter_context(mock.patch.object(client_main, "run_benchmark_batch", return_value=6)) + buf = io.StringIO() + with contextlib.redirect_stdout(buf): + rc = client_main.run_sweep_mode(mode="small", base_args=self.args(queue_dir=queue_dir), + show_end_screen=True, interactive=False, presets_cfg={}) + self.assertEqual(rc, 6) + self.assertIn("nothing was marked finished", buf.getvalue()) + def test_sweep_starts_fresh_when_saved_plan_differs(self) -> None: with tempfile.TemporaryDirectory() as queue_dir: stale = {"encoder": "libx264", "preset": "slow", "crf": 30, diff --git a/client/tests/test_measurement_budget.py b/client/tests/test_measurement_budget.py index 7f5af6d2..d3fa49bb 100644 --- a/client/tests/test_measurement_budget.py +++ b/client/tests/test_measurement_budget.py @@ -108,3 +108,117 @@ def test_gui_and_cli_single_args_preserve_requested_allowance(): assert main.build_arg_parser().parse_args([]).max_duration_minutes == 60 effective=main.build_single_effective_args(base_args=base,encoder='libx264',preset='fast',crf=24) assert effective.max_duration_minutes == 0.25 + + +def test_checkpoint_uploads_terminal_groups_and_retires_accepted_artifacts(tmp_path): + from dataclasses import replace + from test_main_routing import MainRoutingTests, _DummyDashboard + fixture = MainRoutingTests() + clip_a = fixture._quick_clip() + clip_b = replace(clip_a, clip_id="film-grain-1080p24-final", workload_id="film-grain-1080p24-final") + args = fixture._batch_args(str(tmp_path), no_submit=False) + args.local_metrics = False + args.campaign_seed = 23 + args.max_duration_minutes = 1.0 + clock = [0.0] + calls = [] + submissions = [] + def new_budget(minutes, **kwargs): + return campaign.MeasurementBudget(minutes, clock=lambda: clock[0], **kwargs) + def encode(**kwargs): + calls.append(kwargs['artifact_name']) + clock[0] += 10 + campaign.check_measurement_budget() + artifact = Path(kwargs['out_dir']) / kwargs['artifact_name'] + artifact.write_bytes(b'encoded') + return {'artifactPath': str(artifact), 'encoderUsed': 'libx264', 'presetUsed': 'fast', 'fileSizeBytes': 7, + 'encodeStartMonotonicNs': 1_000_000_000, 'encodeEndMonotonicNs': 2_000_000_000, + 'elapsedMs': 1000, 'error': None} + def submit(queue_dir, *, base_url, payload, max_storage_mb, api_key, retries, use_token): + submissions.append(payload) + return "submitted", "run-checkpoint-test", 0 + hardware = main.HardwareInfo('CPU', None, 16, 'OS') + budget_stack = mock.patch.object(main, 'MeasurementBudget', side_effect=new_budget) + with budget_stack, \ + mock.patch.object(main, 'detect_hardware', return_value=main.HardwareInfo('CPU', None, 16, 'OS')), \ + mock.patch.object(main, 'check_compatibility', return_value={}), \ + mock.patch.object(main, 'fetch_baseline_rows', return_value=[]), \ + mock.patch.object(main, '_submit_payload_with_spool', side_effect=submit), \ + mock.patch.object(main, 'ensure_ffmpeg_and_ffprobe', return_value=(True, 'ffmpeg test')), \ + mock.patch.object(main, '_build_protocol_config', + return_value=protocol.ProtocolConfig.for_version('7.1', max_adaptive_repeats=0)), \ + mock.patch.object(main, 'probe_video_stream_metrics', + return_value={'sourceFps': 24, 'sourceDurationSeconds': 5, 'containerFormat': 'mp4'}), \ + mock.patch.object(main, '_probe_artifact_contract', side_effect=lambda path: fixture._artifact_contract()), \ + mock.patch.object(main, '_capture_protocol_environment_snapshot', + return_value=protocol.EnvironmentSnapshot(selected_accelerator='software')), \ + mock.patch.object(main, 'encode_to_artifact', side_effect=encode), \ + mock.patch.object(main, 'BatchRunDashboard', _DummyDashboard): + # Warmups + the first measured round complete; the second measured round + # dies at its last attempt, so one group is terminal and one is partial. + assert main.run_benchmark_batch(hardware=hardware, base_url='https://example.invalid', args=args, + tasks=[{'encoder': 'libx264', 'preset': 'fast', 'crf': 24, 'suiteClip': clip_a}, + {'encoder': 'libx264', 'preset': 'fast', 'crf': 24, 'suiteClip': clip_b}]) == 11 + root = next((tmp_path / 'campaigns').iterdir()) + receipts = sorted(root.glob('submission-*.accepted.json')) + accepted = [json.loads(path.read_text()) for path in receipts] + assert len(receipts) == 2 # Only the terminal group uploads; the partial group waits. + assert all(entry['benchmarkRunId'] == 'run-checkpoint-test' for entry in accepted) + assert len(submissions) == 2 + for entry in accepted: + assert not Path(entry['artifactPath']).exists() # Accepted bytes retire. + # Warmup outputs are never publishable: they release once their attempt + # is durable with a matching hash, leaving verified release evidence. + released = sorted(root.glob('warmup-*.released.json')) + assert len(released) == 2 + assert not [path for path in root.glob('*warmup*.mp4')] + # Every still-unuploaded measured artifact remains byte-complete. + remaining = [path for path in root.glob('*.mp4')] + assert remaining # the partial group's attempts survive + for path in remaining: + assert path.read_bytes() == b'encoded' + status = json.loads(next(root.glob('budget-exhausted-*.json')).read_text()) + assert status['checkpointUploads'] == 2 + # A checkpointed journal reopens without the accepted bytes... + manifest = json.loads((root / 'manifest.json').read_text()) + campaign.CampaignJournal(str(tmp_path), root.name, manifest, 2048) + # ...but a receipt that does not match its attempt fails closed. + receipts[0].write_text(json.dumps({**accepted[0], 'artifactSha256': 'b' * 64})) + with pytest.raises(ValueError, match='Journal artifact missing'): + campaign.CampaignJournal(str(tmp_path), root.name, manifest, 2048) + receipts[0].write_text(json.dumps(accepted[0])) + # Warmup release evidence is validated the same way: tampered evidence + # with the bytes already gone fails the journal reopen closed. + evidence = json.loads(released[0].read_text()) + released[0].write_text(json.dumps({**evidence, 'artifactSha256': 'c' * 64})) + with pytest.raises(ValueError, match='release evidence'): + campaign.CampaignJournal(str(tmp_path), root.name, manifest, 2048) + released[0].write_text(json.dumps(evidence)) + # A corrupt (empty) accepted receipt must never skip the real upload: + # plant one on the unfinished group's pending attempt, bytes present. + accepted_orders = {entry['executionOrder'] for entry in accepted} + pending = [] + for path in sorted(root.glob('attempt-*.json')): + data = json.loads(path.read_text()) + schedule = data.get('schedule', {}) + info = (data.get('metadata') or {}).get('info') or {} + if (schedule.get('phase') == 'measured' + and schedule.get('execution_order') not in accepted_orders + and Path(str(info.get('artifactPath') or '')).is_file()): + pending.append(schedule.get('execution_order')) + assert len(pending) == 1 + (root / f'submission-{pending[0]:06d}.accepted.json').write_text('') + # Completing the campaign uploads the remaining group exactly once more + # and never re-uploads the checkpointed attempts. + clock[0] = 100.0 + assert main.main(['prog', '--resume-campaign', root.name, '--submit', + '--max-duration-minutes', '0.5', '--queue-dir', str(tmp_path)]) == 0 + accepted_files = sorted(root.glob('submission-*.accepted.json')) + assert len(accepted_files) == 4 # two per group, once each + assert len(submissions) == 4 # the corrupt-receipt attempt uploaded exactly once + for path in accepted_files: + entry = json.loads(path.read_text()) + assert entry['benchmarkRunId'] == 'run-checkpoint-test' # real server run id + assert entry['artifactSha256'] != '' and len(entry['artifactSha256']) == 64 + assert not Path(entry['artifactPath']).exists() # all accepted bytes retired + assert (root / 'campaign-complete.json').exists() diff --git a/client/tests/test_owned_artifact_sync.py b/client/tests/test_owned_artifact_sync.py index c4a5ff92..8d744d94 100644 --- a/client/tests/test_owned_artifact_sync.py +++ b/client/tests/test_owned_artifact_sync.py @@ -119,7 +119,7 @@ def test_retained_campaign_traceback_is_opt_in_and_exit_stays_failed(tmp_path, m def failed_lock(): raise OSError(errno.EBADF, 'Bad file descriptor') yield - journal = SimpleNamespace(root=tmp_path, check_budget=lambda: 1000, measurement_lock=failed_lock) + journal = SimpleNamespace(root=tmp_path, records={}, check_budget=lambda: 1000, measurement_lock=failed_lock) with mock.patch.object(main, 'ensure_ffmpeg_and_ffprobe', return_value=(True, 'test')), mock.patch.object(main, '_build_protocol_recipe_specs', return_value=[protocol.RecipeSpec('recipe', protocol.StructuralExpectation())]), mock.patch.object(main, 'CampaignJournal', return_value=journal), mock.patch.object(main, 'physical_source_id', return_value='installation-'+'a'*64), mock.patch('client.identity.runtime_identity', return_value={}), mock.patch.object(main, 'selected_device', return_value={'deviceId':'cpu'}): code = main.run_benchmark_batch(hardware=main.HardwareInfo('CPU',None,16,'OS'), base_url='unused', args=args, tasks=[{'encoder':'libx264','preset':'fast','crf':24,'suiteClip':clip}]) diff --git a/client/tests/test_publication_storage.py b/client/tests/test_publication_storage.py index 14f4ad1b..99521e04 100644 --- a/client/tests/test_publication_storage.py +++ b/client/tests/test_publication_storage.py @@ -23,7 +23,7 @@ def snapshot(root): return {str(p.relative_to(root)): p.read_bytes() for p in root.rglob('*') if p.is_file()} -@pytest.mark.parametrize('retained_folder', ['campaigns/old', 'receipts', 'terminal', 'dead-letter', 'artifacts']) +@pytest.mark.parametrize('retained_folder', ['receipts', 'terminal', 'dead-letter', 'artifacts']) def test_aggregate_counts_all_retained_data_before_copy(tmp_path, retained_folder): queue = tmp_path / 'queue' prior = queue / retained_folder / 'retained.dat' diff --git a/client/tests/test_release_preflight.py b/client/tests/test_release_preflight.py index 91ee0cd1..b7b20448 100644 --- a/client/tests/test_release_preflight.py +++ b/client/tests/test_release_preflight.py @@ -58,7 +58,7 @@ def test_release_json_declares_coherent_frozen_release(self) -> None: payload = json.loads((release_manifest_lib.ROOT_DIR / "release.json").read_text(encoding="utf-8")) self.assertEqual(payload["suiteVersion"], "encodingdb-test-suite-v1") - self.assertEqual(payload["projectVersion"], "1.3.0-rc.4") + self.assertEqual(payload["projectVersion"], "1.3.0-rc.5") self.assertEqual(payload["releaseDate"], "2026-09-22") for tree in ("client", "server"): root = release_manifest_lib.ROOT_DIR / tree / "resources/test_suite_v1" diff --git a/client/tests/test_resume_storage_policy.py b/client/tests/test_resume_storage_policy.py new file mode 100644 index 00000000..72bbd27f --- /dev/null +++ b/client/tests/test_resume_storage_policy.py @@ -0,0 +1,96 @@ +import json +from types import SimpleNamespace +from unittest import mock + +import pytest + +from client import main +from client.campaign import CampaignJournal, atomic_json, journal_path + + +@pytest.mark.parametrize("requested,explicit,expected", [(2048, False, 6144), (512, True, 512), (8192, True, 8192)]) +def test_direct_resume_resolves_policy_and_journal_obeys_it(tmp_path, requested, explicit, expected): + campaign_id = "campaign-1234567890abcdef" + saved = {"seed": 17, "tasks": []} + original = CampaignJournal(str(tmp_path), campaign_id, saved, 6144) + args = SimpleNamespace(queue_dir=str(tmp_path), resume_campaign=campaign_id, + max_storage_mb=requested, max_storage_mb_explicit=explicit, + base_url="https://example.invalid") + + def reopen(**kwargs): + journal = CampaignJournal(str(tmp_path), campaign_id, saved, kwargs["args"].max_storage_mb) + assert journal.max_bytes == expected * 1024 * 1024 + return 0 + + with mock.patch.object(main, "_apply_submission_policy", side_effect=lambda args, **kw: args), \ + mock.patch.object(main, "_preparation_preflight", return_value=0), \ + mock.patch.object(main, "detect_hardware", return_value={}), \ + mock.patch.object(main, "run_benchmark_batch", side_effect=reopen): + assert main._resume_campaign(args) == 0 + assert json.loads((original.root / "budget.json").read_text())["maxStorageMb"] == expected + assert json.loads((original.root / "manifest.json").read_text()) == saved + + +def test_direct_legacy_resume_allows_existing_retained_bytes(tmp_path): + campaign_id = "campaign-1234567890abcdef" + root = journal_path(str(tmp_path), campaign_id) + atomic_json(root / "manifest.json", {"seed": 17, "tasks": []}) + args = SimpleNamespace(queue_dir=str(tmp_path), resume_campaign=campaign_id, + max_storage_mb=2048, max_storage_mb_explicit=False, + base_url="https://example.invalid") + with mock.patch.object(main, "_apply_submission_policy", side_effect=lambda args, **kw: args), \ + mock.patch.object(main, "_preparation_preflight", return_value=0), \ + mock.patch.object(main, "detect_hardware", return_value={}), \ + mock.patch.object(main, "directory_bytes", return_value=2100 * 1024 * 1024), \ + mock.patch.object(main, "run_benchmark_batch", return_value=0) as run: + assert main._resume_campaign(args) == 0 + assert run.call_args.kwargs["args"].max_storage_mb == 3124 + + +def test_resume_prepares_each_clip_once_across_recipes(tmp_path): + campaign_id = "campaign-1234567890abcdef" + task = {"encoder": "libx264", "preset": "fast", "crf": 24, "rateControl": None, "clipId": "clip-a"} + atomic_json(journal_path(str(tmp_path), campaign_id) / "manifest.json", + {"seed": 17, "tasks": [task, dict(task, preset="slow"), dict(task, clipId="clip-b")]}) + args = SimpleNamespace(queue_dir=str(tmp_path), resume_campaign=campaign_id, + max_storage_mb=2048, max_storage_mb_explicit=False, base_url="https://example.invalid") + with mock.patch.object(main, "_apply_submission_policy", side_effect=lambda args, **kw: args), \ + mock.patch.object(main, "_preparation_preflight", return_value=0), \ + mock.patch.object(main, "detect_hardware", return_value={}), \ + mock.patch.object(main, "_prepare_named_suite_clip", side_effect=lambda clip: object()) as prepare, \ + mock.patch.object(main, "run_benchmark_batch", return_value=0) as run: + assert main._resume_campaign(args) == 0 + assert [call.args[0] for call in prepare.call_args_list] == ["clip-a", "clip-b"] + tasks = run.call_args.kwargs["tasks"] + assert tasks[0]["suiteClip"] is tasks[1]["suiteClip"] + + +@pytest.mark.parametrize("explicit_duration,calls", [(False, 2), (True, 1)]) +def test_saved_sweep_restores_plan_attempt_cap_and_continues_checkpoints(tmp_path, explicit_duration, calls): + campaign_id = "campaign-1234567890abcdef" + tasks = [{"encoder": "libx264", "preset": str(i), "crf": 24, "rateControl": None, "clipId": "clip-a"} + for i in range(40)] + atomic_json(journal_path(str(tmp_path), campaign_id) / "manifest.json", + {"seed": 17, "sweepMode": "large", "tasks": tasks}) + argv = ["--cli", "--resume-campaign", campaign_id, "--queue-dir", str(tmp_path)] + if explicit_duration: + argv += ["--max-duration-minutes", "1"] + args = main.build_arg_parser().parse_args(argv) + invoked = [] + + def batch(**kwargs): + invoked.append(kwargs) + main.config._BATCH_ATTEMPTS_RECORDED = 1 + return 11 if len(invoked) == 1 else 0 + + with mock.patch.object(main, "_apply_submission_policy", side_effect=lambda args, **kw: args), \ + mock.patch.object(main, "_preparation_preflight", return_value=0), \ + mock.patch.object(main, "detect_hardware", return_value={}), \ + mock.patch.object(main, "_prepare_named_suite_clip", return_value=object()) as prepare, \ + mock.patch.object(main, "run_benchmark_batch", side_effect=batch): + assert main._resume_campaign(args) == (11 if explicit_duration else 0) + assert len(invoked) == calls + assert prepare.call_count == 1 + assert args.max_storage_mb == 6144 + protocol = main._build_protocol_config() + assert args.max_attempts == 40 * (protocol.warmup_runs + protocol.minimum_measured_runs + protocol.max_adaptive_repeats) diff --git a/client/tests/test_source_preparation_dedup.py b/client/tests/test_source_preparation_dedup.py index 50b6033a..6827084d 100644 --- a/client/tests/test_source_preparation_dedup.py +++ b/client/tests/test_source_preparation_dedup.py @@ -139,3 +139,63 @@ def test_source_contract_preserves_existing_metrics_fallback_without_extra_probe assert recipe.expectation.duration_s == 1 / 12 assert recipe.expectation.avg_frame_rate == 24 assert recipe.expectation.frame_count == 2 + + +def _two_clip_prepared(manifest): + return [suite.ensure_suite_clip(manifest.clips[0]), suite.ensure_suite_clip(manifest.clips[1])] + + +def test_recipe_source_contract_probes_each_distinct_clip_once(): + from test_suite_v1 import small_media_fixture + with small_media_fixture() as (_, manifest): + clips = _two_clip_prepared(manifest) + main._SOURCE_CONTRACT_CACHE.clear() + tasks = [{'encoder': 'libx264', 'preset': 'fast', 'crf': 24, 'suiteClip': clips[index % 2]} + for index in range(8)] + with mock.patch.object(main, '_probe_artifact_contract', + wraps=main._probe_artifact_contract) as contract: + recipes = main._build_protocol_recipe_specs( + tasks, default_input_path=clips[0].path, default_input_hash=clips[0].input_hash) + assert contract.call_count == 2 # Once per distinct frozen clip, never once per recipe. + assert len(recipes) == 8 + assert all(recipe.expectation.width == 32 for recipe in recipes) + + +def test_recipe_source_contract_reports_advancing_progress(): + from test_suite_v1 import small_media_fixture + with small_media_fixture() as (_, manifest): + clips = [suite.ensure_suite_clip(clip) for clip in manifest.clips[:3]] + main._SOURCE_CONTRACT_CACHE.clear() + events = [] + tasks = [{'encoder': 'libx264', 'preset': 'fast', 'crf': 24, 'suiteClip': clips[index % 3]} + for index in range(9)] + with campaign.PreparationScope(progress=lambda stage, **details: events.append((stage, details))).activate(): + main._build_protocol_recipe_specs( + tasks, default_input_path=clips[0].path, default_input_hash=clips[0].input_hash) + contracts = [details for stage, details in events if stage == 'source-contract'] + assert [details['completed'] for details in contracts] == [1, 2, 3] + assert all(details['total'] == 3 for details in contracts) + assert {details['clipId'] for details in contracts} == {clip.clip_id for clip in clips} + + +def test_recipe_source_contract_cancellation_stops_before_next_probe(): + from test_suite_v1 import small_media_fixture + with small_media_fixture() as (_, manifest): + clips = _two_clip_prepared(manifest) + main._SOURCE_CONTRACT_CACHE.clear() + stop = threading.Event() + real_probe = main._probe_artifact_contract + probed = [] + def probe_once(path): + probed.append(path) + result = real_probe(path) + stop.set() # Cancel lands after the first probe; the next task's fence must fire. + return result + with mock.patch.object(main, '_probe_artifact_contract', side_effect=probe_once): + tasks = [{'encoder': 'libx264', 'preset': 'fast', 'crf': 24, 'suiteClip': clips[index % 2]} + for index in range(8)] + with campaign.PreparationScope(stop).activate(): + with pytest.raises(KeyboardInterrupt): + main._build_protocol_recipe_specs( + tasks, default_input_path=clips[0].path, default_input_hash=clips[0].input_hash) + assert len(probed) == 1 # The cancel fence halts the plan before the next full decode. diff --git a/client/tests/test_spool.py b/client/tests/test_spool.py index 70fcb21c..ecc60258 100644 --- a/client/tests/test_spool.py +++ b/client/tests/test_spool.py @@ -293,6 +293,52 @@ def test_inspect_and_cleanup_spool_preserve_pending_entries(self) -> None: server.server_close() thread.join(timeout=2) + def test_unrelated_retained_campaign_does_not_block_this_upload(self) -> None: + # Review scenario: a Large campaign retaining >2 GiB must not consume the + # allowance of a Small run spooling its own measured artifact. + with tempfile.TemporaryDirectory() as queue_dir, tempfile.TemporaryDirectory() as source_dir: + other_root = os.path.join(queue_dir, "campaigns", "campaign-" + "f" * 16) + os.makedirs(other_root) + with open(os.path.join(other_root, "retained.bin"), "wb") as bulk: + bulk.truncate(2100 * 1024 * 1024) + source_path = os.path.join(queue_dir, "campaigns", "campaign-" + "a" * 16, "artifact.mp4") + os.makedirs(os.path.dirname(source_path), exist_ok=True) + with open(source_path, "wb") as handle: + handle.write(b"test") + payload = self._authoritative_payload(source_path) + path, _entry = spool_payload(queue_dir, payload, max_storage_mb=2048) + managed_path = load_spool_entry(path)["payload"]["artifactPath"] + self.assertTrue(os.path.exists(managed_path)) + with mock.patch("client.spool.submit_artifact_submission", + return_value={"analyses": [{"vmafMean": 95.25}]}): + stats = replay_spool(queue_dir, base_url="http://127.0.0.1:9", api_key="", + retries=1, use_token=False) + self.assertEqual(stats.submitted, 1) + self.assertEqual(count_pending_entries(queue_dir), 0) + self.assertFalse(os.path.exists(managed_path)) # accepted upload retires staging + self.assertTrue(os.path.exists(os.path.join(other_root, "retained.bin"))) + + def test_spool_capacity_charges_own_campaign_and_honors_free_floor(self) -> None: + from types import SimpleNamespace + with tempfile.TemporaryDirectory() as queue_dir, tempfile.TemporaryDirectory() as source_dir: + mine = os.path.join(queue_dir, "campaigns", "campaign-" + "a" * 16) + os.makedirs(mine) + bulk = os.path.join(mine, "retained.bin") + with open(bulk, "wb") as handle: + handle.truncate(2100 * 1024 * 1024) + source_path = os.path.join(queue_dir, "campaigns", "campaign-" + "a" * 16, "artifact.mp4") + os.makedirs(os.path.dirname(source_path), exist_ok=True) + with open(source_path, "wb") as handle: + handle.write(b"test") + payload = self._authoritative_payload(source_path) + with self.assertRaisesRegex(OSError, "this campaign's retained attempts"): + spool_payload(queue_dir, payload, max_storage_mb=2048) + os.remove(bulk) # own campaign now small; disk floor must still protect the volume + with mock.patch("client.spool.shutil.disk_usage", + return_value=SimpleNamespace(total=1024 ** 4, used=0, free=500 * 1024 * 1024)): + with self.assertRaisesRegex(OSError, "safety"): + spool_payload(queue_dir, payload, max_storage_mb=2048) + if __name__ == "__main__": unittest.main() diff --git a/client/tests/test_windows_gui.py b/client/tests/test_windows_gui.py index 2c85de48..5178a33b 100644 --- a/client/tests/test_windows_gui.py +++ b/client/tests/test_windows_gui.py @@ -119,6 +119,11 @@ def test_initialized_controls_and_keyboard_handlers_use_real_running_guards(self self.assertEqual((app._selected_encoder(), app._selected_preset(), app.crf_var.get()), ("libx264", "fast", 0)) app._handle_event({"type": "preparation_progress", "stage": "probe", "path": "test.mkv"}) self.assertEqual(app.stage_var.get(), "Preparing: probe") + app._handle_event({"type": "preparation_progress", "stage": "source-contract", + "clipId": "talking-head-1080p24-final", "completed": 3, "total": 7}) + self.assertEqual(app.stage_var.get(), "Preparing: source-contract") + self.assertIn("(3/7)", app.summary_var.get()) + self.assertIn("talking-head-1080p24-final", app.summary_var.get()) self.assertIn("Stop is available", app.summary_var.get()) app._handle_event({"type": "preparation_progress", "stage": "recovery", "path": "recovered", "message": "installed a fully verified copy"}) @@ -397,7 +402,7 @@ def test_explicit_allowance_reported_when_set(self): app, _root, bindings, _tk = self.build() log_lines = [] app.base_args.max_duration_minutes = 15.0 - app.base_args.explicit_max_duration_minutes = True + app.base_args.max_duration_minutes_explicit = True with mock.patch.object(gui.threading, "Thread", FakeThread), \ mock.patch.object(app, "_append_log", log_lines.append): bindings[""](None) diff --git a/client/ui.py b/client/ui.py index 29d1d636..5335c75d 100644 --- a/client/ui.py +++ b/client/ui.py @@ -287,17 +287,31 @@ def _format_duration(seconds: float) -> str: return " ".join(parts) -def print_end_screen(completed_count: int, elapsed_seconds: float) -> None: +def print_end_screen(completed_count: int, elapsed_seconds: float, status: str = "complete", + recovery: Optional[str] = None) -> None: + """Render the terminal state honestly; only rc-0 work may claim completion.""" time_str = _format_duration(elapsed_seconds) + headline, title, plain = { + "complete": ("[ok] Benchmark run complete [/ok]", "Thank You", "Benchmark complete."), + "paused": ("[accent2] Campaign checkpointed — retained, not finished [/accent2]", "Not Finished", + "Campaign checkpointed and retained; not finished."), + "interrupted": ("[accent2] Run interrupted — retained, not finished [/accent2]", "Interrupted", + "Run interrupted; progress retained; not finished."), + "failed": ("[accent] Run did not complete — nothing was marked finished [/accent]", "Failed", + "Run did not complete; nothing was marked finished."), + }.get(status, ("[accent] Run ended with an unknown state [/accent]", "Failed", + "Run ended with an unknown state.")) + recovery_line = f"\n[muted]{recovery}[/muted]" if recovery else "" if _rich_tty(): body = ( - "[ok] Benchmark run complete [/ok]\n\n" + f"{headline}\n\n" f"[muted]Submitted data points:[/muted] [accent]{completed_count}[/accent]\n" - f"[muted]Time donated:[/muted] [accent2]{time_str}[/accent2]" + f"[muted]Time donated:[/muted] [accent2]{time_str}[/accent2]{recovery_line}" ) - _console.print(_Panel(body, title="[title] Thank You [/title]", border_style="accent2")) + _console.print(_Panel(body, title=f"[title] {title} [/title]", border_style="accent2")) else: - print(f"Benchmark complete. Submitted {completed_count} data points in {time_str}.") + print(f"{plain} Submitted {completed_count} data points in {time_str}." + + (f" {recovery}" if recovery else "")) def print_benchmark_result(payload: Dict[str, Any], relative_file_size_pct: Optional[float]) -> None: diff --git a/client/windows_gui.py b/client/windows_gui.py index 2db1c275..2a2c431a 100644 --- a/client/windows_gui.py +++ b/client/windows_gui.py @@ -369,6 +369,17 @@ def _stop_shortcut(self, _event: Any = None) -> str: def _start_run(self) -> None: if self.running or self._upload_active(): return + try: # Advisory only; run_benchmark_batch refuses authoritatively before any preparation. + active = client_main.active_collection(str(self.base_args.queue_dir)) + except Exception: + active = None + if active is not None: + who = f" (campaign {active['campaignId']}, PID {active['pid']})" if active.get("campaignId") else "" + self.summary_var.set(f"Another collection is actively running in this queue{who}. " + f"Its checkpoints continue automatically - let it finish, or " + f"stop/cancel that run first. Retry Queued Uploads stays available.") + self._append_log("Start refused: active collection detected") + return mode = self.mode_var.get().strip() mode_key = GUI_MODE_BY_LABEL.get(mode) if mode_key is None and ( @@ -400,7 +411,7 @@ def _start_run(self) -> None: run_args.pause_on_exit = False run_args.menu = False self._browse_shown = False - if getattr(run_args, "explicit_max_duration_minutes", False): + if getattr(run_args, "max_duration_minutes_explicit", False): self._append_log( f"Explicit measurement allowance: {float(run_args.max_duration_minutes):g} minutes; " "the run stops there with the campaign saved for a later continuation." @@ -518,7 +529,12 @@ def _handle_event(self, event: Dict[str, Any]) -> None: return label = event.get("clipId") or os.path.basename(str(event.get("path") or "")) done, total = event.get("completedBytes"), event.get("totalBytes") - amount = f" ({done}/{total} bytes)" if done is not None and total else "" + if done is not None and total: + amount = f" ({done}/{total} bytes)" + elif event.get("completed") is not None and event.get("total"): + amount = f" ({event['completed']}/{event['total']})" + else: + amount = "" self.summary_var.set(f"Preparing {label}{amount}; Stop is available") return diff --git a/frontend/app/run/page.tsx b/frontend/app/run/page.tsx index 92d95411..62caf52b 100644 --- a/frontend/app/run/page.tsx +++ b/frontend/app/run/page.tsx @@ -67,7 +67,7 @@ export default function RunPage() {

Windows: SmartScreen warns because the executable is unsigned. Verify the SHA-256 above first, and continue only if you trust the source; the page gives no bypass tool or automation.

Superseded packaged builds ({supersededTag})

-

The {supersededTag} Windows build names every failure cause but stops at a cache folder protected against the current user until that folder is deleted with administrator rights; {projectTag} recovers automatically instead, with no manual repair. Its macOS and Linux binaries are byte-identical to the current ones.

+

The {supersededTag} macOS build predates the preparation, retained-campaign storage, and resume repairs in {projectTag}. The verified Windows and Linux binaries are unchanged in this release.