Files

583 lines
19 KiB
Python

#!/usr/bin/env python3
"""Drive selective, negative, cancellation, and recovery E2E jobs."""
from __future__ import annotations
import argparse
import json
import sys
import time
import urllib.error
import urllib.parse
import urllib.request
CONTROL = "http://127.0.0.1:18081/test/v1"
class HttpError(RuntimeError):
def __init__(self, status: int, body: object):
super().__init__(f"HTTP {status}: {body}")
self.status = status
self.body = body
def get_json(url: str):
try:
with urllib.request.urlopen(url, timeout=35) as response:
return json.load(response)
except urllib.error.HTTPError as exc:
raise HttpError(exc.code, _error_body(exc)) from exc
def post_json(url: str, value: object):
request = urllib.request.Request(
url,
data=json.dumps(value, separators=(",", ":")).encode(),
method="POST",
headers={"Content-Type": "application/json"},
)
try:
with urllib.request.urlopen(request, timeout=35) as response:
return json.load(response)
except urllib.error.HTTPError as exc:
raise HttpError(exc.code, _error_body(exc)) from exc
def _error_body(error: urllib.error.HTTPError):
try:
return json.load(error)
except (ValueError, OSError):
return error.read().decode("utf-8", "replace")
def indices(value: str) -> list[int]:
if not value:
return []
result = [int(item) for item in value.split(",")]
if result != sorted(set(result)) or any(item < 0 for item in result):
raise ValueError("indices must be unique, sorted, non-negative integers")
return result
def selection_indices(value: dict[str, object]) -> list[int]:
result = []
for item in value.get("ranges", []):
first = int(item.get("first", 0))
last = int(item.get("last", 0))
result.extend(range(first, last + 1))
return result
def resource_id(info_hash: str) -> dict[str, str]:
key = "info_hash_v1_hex" if len(info_hash) == 40 else "info_hash_v2_hex"
return {key: info_hash}
def qb_json(endpoint: str, path: str, **parameters: str):
query = urllib.parse.urlencode(parameters)
suffix = f"?{query}" if query else ""
return get_json(f"{endpoint}{path}{suffix}")
def transfer_preview(args: argparse.Namespace):
request = {
"operation": args.operation,
"source_client_id": args.source_client,
"target_client_id": args.target_client,
"resource_id": resource_id(args.info_hash),
}
selected = indices(args.selected)
if selected:
request["selected_file_indices"] = selected
return post_json(f"{CONTROL}/jobs/preview", request)
def create_previewed(preview: dict[str, object]) -> str:
definition = preview["definition"]
job_id = definition["job_id"]
created = post_json(f"{CONTROL}/jobs", {
"preview_revision": preview["preview_revision"],
"definition": definition,
})
if created["job_id"] != job_id:
raise RuntimeError("control returned the wrong job")
return job_id
def wait_job(job_id: str, expected: str, timeout: float = 180):
deadline = time.monotonic() + timeout
last = None
while time.monotonic() < deadline:
last = get_json(f"{CONTROL}/jobs/{job_id}")
state = last["state"]
if state == expected:
return last
if state in {
"JOB_STATE_FAILED",
"JOB_STATE_CANCELLED",
"JOB_STATE_SUCCEEDED",
}:
raise RuntimeError(
f"job reached {state} while waiting for {expected}: {last}"
)
time.sleep(0.25)
raise RuntimeError(
f"timed out waiting for {job_id} to reach {expected}: {last}"
)
def wait_latest(
job_id: str,
*,
event_type: str,
step: str | None = None,
minimum_sequence: int = 0,
timeout: float = 90,
):
deadline = time.monotonic() + timeout
last = None
while time.monotonic() < deadline:
last = get_json(f"{CONTROL}/jobs/{job_id}")
if last["state"] == "JOB_STATE_FAILED":
raise RuntimeError(f"job failed while waiting for event: {last}")
event = last.get("latest_event") or {}
progress = event.get("progress") or {}
if (
event.get("type") == event_type
and int(event.get("sequence", 0)) >= minimum_sequence
and (step is None or progress.get("step") == step)
):
return last
time.sleep(0.1)
raise RuntimeError(f"timed out waiting for job event: {last}")
def advance_for(job_id: str, action: str):
outcomes = post_json(f"{CONTROL}/scheduler/advance", {})
matching = [
item for item in outcomes
if item.get("job_id") == job_id and item.get("action") == action
]
if len(matching) != 1:
raise RuntimeError(
f"expected scheduler action {action} for {job_id}: {outcomes}"
)
return matching[0]
def assert_transfer_preview(
preview: dict[str, object], expected: list[int], delta: list[int]
) -> None:
transfer = preview["definition"]["transfer"]
actual_requested = selection_indices(transfer["requested_files"])
actual_delta = selection_indices(transfer["transfer_delta_files"])
if actual_requested != expected or actual_delta != delta:
raise RuntimeError(
"unexpected transfer selection: "
f"requested={actual_requested} delta={actual_delta}"
)
def command_transfer(args: argparse.Namespace) -> None:
preview = transfer_preview(args)
expected = indices(args.expected_selection or args.selected)
delta = indices(args.expected_delta or args.expected_selection or args.selected)
if expected:
assert_transfer_preview(preview, expected, delta)
job_id = create_previewed(preview)
job = wait_job(job_id, "JOB_STATE_SUCCEEDED")
if not job["committed"]:
raise RuntimeError("successful transfer was not committed")
print(json.dumps({
"job_id": job_id,
"operation": args.operation,
"selection": expected,
"delta": delta,
"state": job["state"],
}, sort_keys=True))
def command_preview_error(args: argparse.Namespace) -> None:
try:
if args.operation == "evict_cache":
post_json(f"{CONTROL}/jobs/preview", {
"operation": args.operation,
"cache_client_id": args.cache_client,
"resource_id": resource_id(args.info_hash),
})
else:
transfer_preview(args)
except HttpError as exc:
rendered = json.dumps(exc.body, sort_keys=True)
if exc.status != args.status or args.contains not in rendered:
raise RuntimeError(
f"unexpected preview failure: status={exc.status} body={rendered}"
) from exc
print(json.dumps({
"expected_status": exc.status,
"matched": args.contains,
}, sort_keys=True))
return
raise RuntimeError("preview unexpectedly succeeded")
def command_stale_preview(args: argparse.Namespace) -> None:
preview = transfer_preview(args)
definition = preview["definition"]
definition["resource_display_name"] += "-tampered"
try:
post_json(f"{CONTROL}/jobs", {
"preview_revision": preview["preview_revision"],
"definition": definition,
})
except HttpError as exc:
if exc.status != 409:
raise
print(json.dumps({"stale_preview_rejected": True}, sort_keys=True))
return
raise RuntimeError("tampered preview revision was accepted")
def command_evict(args: argparse.Namespace) -> None:
preview = post_json(f"{CONTROL}/jobs/preview", {
"operation": "evict_cache",
"cache_client_id": args.cache_client,
"resource_id": resource_id(args.info_hash),
})
coverage = preview["definition"]["eviction"]["archive_coverage"]
if len(coverage) < args.minimum_coverage_proofs:
raise RuntimeError(
f"expected at least {args.minimum_coverage_proofs} coverage proofs"
)
job_id = create_previewed(preview)
job = wait_job(job_id, "JOB_STATE_SUCCEEDED")
if not job["committed"]:
raise RuntimeError("successful eviction was not committed")
print(json.dumps({
"job_id": job_id,
"coverage_proofs": len(coverage),
"state": job["state"],
}, sort_keys=True))
_STEPS = [
(
"source_stage",
"source_stage",
"JOB_STEP_KIND_SOURCE_STAGE",
),
(
"syncthing_transfer",
"syncthing_transfer",
"JOB_STEP_KIND_SYNCTHING_TRANSFER",
),
(
"target_materialize",
"target_materialize",
"JOB_STEP_KIND_TARGET_MATERIALIZE",
),
]
def command_drive(args: argparse.Namespace) -> None:
post_json(f"{CONTROL}/scheduler/pause", {})
preview = transfer_preview(args)
expected = indices(args.expected_selection or args.selected)
delta = indices(args.expected_delta or args.expected_selection or args.selected)
if expected:
assert_transfer_preview(preview, expected, delta)
job_id = create_previewed(preview)
advance_for(job_id, "assign_source")
wait_latest(
job_id, event_type="JOB_EVENT_TYPE_ASSIGNED", minimum_sequence=1
)
advance_for(job_id, "assign_target")
wait_latest(
job_id, event_type="JOB_EVENT_TYPE_ASSIGNED", minimum_sequence=2
)
if args.through == "assigned":
print(job_id)
return
for boundary, action, step in _STEPS:
advance_for(job_id, action)
wait_latest(
job_id,
event_type="JOB_EVENT_TYPE_STEP_SUCCEEDED",
step=step,
)
if args.through == boundary:
print(job_id)
return
raise RuntimeError(f"unsupported drive boundary: {args.through}")
def command_create_paused(args: argparse.Namespace) -> None:
post_json(f"{CONTROL}/scheduler/pause", {})
preview = transfer_preview(args)
expected = indices(args.expected_selection or args.selected)
delta = indices(args.expected_delta or args.expected_selection or args.selected)
if expected:
assert_transfer_preview(preview, expected, delta)
print(create_previewed(preview))
def command_queued_cancel(args: argparse.Namespace) -> None:
post_json(f"{CONTROL}/scheduler/pause", {})
job_id = create_previewed(transfer_preview(args))
cancelled = post_json(f"{CONTROL}/jobs/{job_id}/cancel", {})
if cancelled.get("disposition") != "removed":
raise RuntimeError(f"queued cancellation was not record-only: {cancelled}")
try:
get_json(f"{CONTROL}/jobs/{job_id}")
except HttpError as exc:
if exc.status == 404:
print(json.dumps({
"job_id": job_id,
"record_removed": True,
}, sort_keys=True))
return
raise
raise RuntimeError("queued cancelled job record still exists")
def command_cancel(args: argparse.Namespace) -> None:
post_json(f"{CONTROL}/jobs/{args.job_id}/cancel", {})
deadline = time.monotonic() + 90
last = None
while time.monotonic() < deadline:
last = get_json(f"{CONTROL}/jobs/{args.job_id}")
if last["state"] == "JOB_STATE_CANCELLED":
if last["committed"]:
raise RuntimeError("cancelled precommit job became committed")
print(json.dumps({
"job_id": args.job_id,
"state": last["state"],
"committed": last["committed"],
}, sort_keys=True))
return
if last["state"] == "JOB_STATE_FAILED":
raise RuntimeError(f"cancellation failed: {last}")
post_json(f"{CONTROL}/scheduler/advance", {})
time.sleep(0.2)
raise RuntimeError(f"timed out waiting for cancellation: {last}")
def command_advance(args: argparse.Namespace) -> None:
outcome = advance_for(args.job_id, args.action)
print(json.dumps(outcome, sort_keys=True))
def command_resume(args: argparse.Namespace) -> None:
post_json(f"{CONTROL}/scheduler/resume", {})
job = wait_job(args.job_id, "JOB_STATE_SUCCEEDED")
if not job["committed"]:
raise RuntimeError("resumed job was not committed")
print(json.dumps({
"job_id": args.job_id,
"state": job["state"],
"committed": job["committed"],
}, sort_keys=True))
def command_expect_failure(args: argparse.Namespace) -> None:
post_json(f"{CONTROL}/scheduler/resume", {})
deadline = time.monotonic() + 120
last = None
while time.monotonic() < deadline:
last = get_json(f"{CONTROL}/jobs/{args.job_id}")
if last["state"] == "JOB_STATE_FAILED":
message = (
(last.get("latest_event") or {})
.get("error", {})
.get("message", "")
)
if args.contains not in message:
raise RuntimeError(
f"failure reason did not contain {args.contains!r}: {last}"
)
if last["committed"]:
raise RuntimeError("hostile precondition failure committed")
print(json.dumps({
"job_id": args.job_id,
"state": last["state"],
"matched": args.contains,
}, sort_keys=True))
return
if last["state"] in {
"JOB_STATE_SUCCEEDED",
"JOB_STATE_CANCELLED",
}:
raise RuntimeError(f"hostile job reached wrong terminal state: {last}")
time.sleep(0.25)
raise RuntimeError(f"timed out waiting for hostile failure: {last}")
def command_assert_qb(args: argparse.Namespace) -> None:
expected = indices(args.selected)
deadline = time.monotonic() + args.timeout
last = None
while time.monotonic() < deadline:
records = qb_json(
args.endpoint, "/torrents/info", hashes=args.info_hash
)
if args.absent:
if records == []:
print(json.dumps({"absent": True}, sort_keys=True))
return
elif len(records) == 1:
runtime_state = str(records[0].get("state", ""))
files = qb_json(
args.endpoint, "/torrents/files", hash=args.info_hash
)
selected = [
int(item["index"]) for item in files
if int(item.get("priority", 0)) > 0
]
selected_complete = [
int(item["index"]) for item in files
if (
int(item.get("priority", 0)) > 0
and float(item.get("progress", 0)) >= 1
)
]
last = {
"runtime_state": runtime_state,
"selected": selected,
"selected_complete": selected_complete,
}
no_downloaded = True
if args.no_downloaded:
properties = qb_json(
args.endpoint,
"/torrents/properties",
hash=args.info_hash,
)
downloaded = int(properties.get("total_downloaded", -1))
last["total_downloaded"] = downloaded
no_downloaded = downloaded == 0
if (
selected == expected
and selected_complete == expected
and runtime_state
not in {"checkingUP", "checkingDL", "checkingResumeData"}
and no_downloaded
):
print(json.dumps(last, sort_keys=True))
return
else:
last = records
time.sleep(0.25)
raise RuntimeError(f"qBittorrent assertion timed out: {last}")
def add_transfer_arguments(parser: argparse.ArgumentParser) -> None:
parser.add_argument(
"--operation", choices=("archive", "unarchive"), required=True
)
parser.add_argument("--source-client", required=True)
parser.add_argument("--target-client", required=True)
parser.add_argument("--info-hash", required=True)
parser.add_argument("--selected", default="")
parser.add_argument("--expected-selection", default="")
parser.add_argument("--expected-delta", default="")
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser()
commands = parser.add_subparsers(dest="command", required=True)
transfer = commands.add_parser("transfer")
add_transfer_arguments(transfer)
transfer.set_defaults(run=command_transfer)
preview_error = commands.add_parser("preview-error")
preview_error.add_argument(
"--operation",
choices=("archive", "unarchive", "evict_cache"),
required=True,
)
preview_error.add_argument("--source-client", default="")
preview_error.add_argument("--target-client", default="")
preview_error.add_argument("--cache-client", default="")
preview_error.add_argument("--info-hash", required=True)
preview_error.add_argument("--selected", default="")
preview_error.add_argument("--status", type=int, default=409)
preview_error.add_argument("--contains", required=True)
preview_error.set_defaults(run=command_preview_error)
stale = commands.add_parser("stale-preview")
add_transfer_arguments(stale)
stale.set_defaults(run=command_stale_preview)
evict = commands.add_parser("evict")
evict.add_argument("--cache-client", required=True)
evict.add_argument("--info-hash", required=True)
evict.add_argument("--minimum-coverage-proofs", type=int, default=1)
evict.set_defaults(run=command_evict)
drive = commands.add_parser("drive")
add_transfer_arguments(drive)
drive.add_argument(
"--through",
choices=(
"assigned",
"source_stage",
"syncthing_transfer",
"target_materialize",
),
required=True,
)
drive.set_defaults(run=command_drive)
create_paused = commands.add_parser("create-paused")
add_transfer_arguments(create_paused)
create_paused.set_defaults(run=command_create_paused)
queued_cancel = commands.add_parser("queued-cancel")
add_transfer_arguments(queued_cancel)
queued_cancel.set_defaults(run=command_queued_cancel)
cancel = commands.add_parser("cancel")
cancel.add_argument("--job-id", required=True)
cancel.set_defaults(run=command_cancel)
advance = commands.add_parser("advance")
advance.add_argument("--job-id", required=True)
advance.add_argument("--action", required=True)
advance.set_defaults(run=command_advance)
resume = commands.add_parser("resume")
resume.add_argument("--job-id", required=True)
resume.set_defaults(run=command_resume)
failure = commands.add_parser("expect-failure")
failure.add_argument("--job-id", required=True)
failure.add_argument("--contains", required=True)
failure.set_defaults(run=command_expect_failure)
qb = commands.add_parser("assert-qb")
qb.add_argument("--endpoint", required=True)
qb.add_argument("--info-hash", required=True)
qb.add_argument("--selected", default="")
qb.add_argument("--absent", action="store_true")
qb.add_argument("--no-downloaded", action="store_true")
qb.add_argument("--timeout", type=float, default=60)
qb.set_defaults(run=command_assert_qb)
return parser
def main() -> int:
args = build_parser().parse_args()
args.run(args)
return 0
if __name__ == "__main__":
try:
raise SystemExit(main())
except Exception as exc:
print(f"complex E2E scenario failed: {exc}", file=sys.stderr)
raise