Download scripts/live_e2e.py from User1342/distinct: direct link, hf CLI and curl.
- Browser
- Download file 21.8 kB
-
https://huggingface.co/spaces/User1342/distinct/resolve/main/scripts/live_e2e.py
- Command line
-
hf download hf://spaces/User1342/distinct/scripts/live_e2e.py
-
curl -L -o live_e2e.py https://huggingface.co/spaces/User1342/distinct/resolve/main/scripts/live_e2e.py
21.8 kB
| """Run the whole product, drive it with a browser, and photograph it. | |
| **What this proves and what it does not.** It starts the real server, a real | |
| worker and a real browser, signs two different people in through the real OAuth | |
| code, redeems a real access code and submits a real run, then saves what the | |
| browser saw. Two things are stood in for, and both are named on the screen: | |
| * **the identity provider.** Hugging Face will not answer an OAuth app that | |
| does not exist, and registering one needs somebody's account and a public | |
| redirect address. So a provider speaking their documented flow runs on | |
| loopback instead. The client code under test is untouched: the same | |
| discovery, the same PKCE challenge, the same nonce, the same HTTP Basic | |
| authentication at the token endpoint. What this cannot show is that | |
| huggingface.co behaves as its documentation says. Its live discovery | |
| document has been read and matches; that is a weaker claim and it is the | |
| one being made. | |
| * **the model.** `--demo-runner` answers deterministically rather than loading | |
| gigabytes of weights. The pairing, the signed snapshots, the queue, the | |
| offer, the acceptance, the result and the energy accounting are the real | |
| ones; only the sentence at the end is not a model's. | |
| **What it demonstrates, which is the thing worth demonstrating.** Alice signs | |
| in, generates a pairing code, and starts a worker. Bob signs in and sees | |
| nothing at all, because signing in is necessary and not sufficient. Bob pastes | |
| the access code the worker printed, and one worker appears. Bob submits a | |
| request and it runs on Alice's machine. Every one of those is a screenshot, | |
| and the third is the one that matters. | |
| """ | |
| from __future__ import annotations | |
| import argparse | |
| import base64 | |
| import hashlib | |
| import json | |
| import os | |
| import re | |
| import shutil | |
| import socket | |
| import subprocess | |
| import sys | |
| import threading | |
| import time | |
| import urllib.parse | |
| from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer | |
| from pathlib import Path | |
| REPO_ROOT = Path(__file__).resolve().parent.parent | |
| sys.path.insert(0, str(REPO_ROOT)) | |
| CODE = "the-authorisation-code" | |
| TOKEN = "the-access-token" | |
| CLIENT_ID = "distinct-live-e2e" | |
| CLIENT_SECRET = "a-secret-only-these-two-processes-know" | |
| #: Who the stand-in provider says is signing in. Changed between the two | |
| #: browser contexts, which is how one provider serves two people. | |
| # Enough people to test the thing that matters about an access code: that it | |
| # is not a one-to-one handshake. Two names could only ever demonstrate "the | |
| # owner and someone else", which is the case that would also pass if the code | |
| # were consumed on first use. | |
| PEOPLE = { | |
| "alice": {"sub": "hf-alice-0001", "preferred_username": "alice", "name": "Alice"}, | |
| "bob": {"sub": "hf-bob-0002", "preferred_username": "bob", "name": "Bob"}, | |
| "carol": {"sub": "hf-carol-0003", "preferred_username": "carol", "name": "Carol"}, | |
| "dave": {"sub": "hf-dave-0004", "preferred_username": "dave", "name": "Dave"}, | |
| "eve": {"sub": "hf-eve-0005", "preferred_username": "eve", "name": "Eve"}, | |
| } | |
| def _id_token(claims: dict) -> str: | |
| def part(payload: dict) -> str: | |
| return base64.urlsafe_b64encode(json.dumps(payload).encode()).decode().rstrip("=") | |
| return f"{part({'alg': 'none'})}.{part(claims)}.signature" | |
| class Provider(BaseHTTPRequestHandler): | |
| """Hugging Face's documented OAuth flow, on loopback, for one run.""" | |
| protocol_version = "HTTP/1.1" | |
| seen: dict = {"who": "alice"} | |
| def log_message(self, *args) -> None: # noqa: A003 - framework name | |
| return None | |
| def _json(self, payload: dict, status: int = 200) -> None: | |
| body = json.dumps(payload).encode() | |
| self.send_response(status) | |
| self.send_header("Content-Type", "application/json") | |
| self.send_header("Content-Length", str(len(body))) | |
| self.end_headers() | |
| self.wfile.write(body) | |
| def do_GET(self) -> None: # noqa: N802 - framework name | |
| path, _, query = self.path.partition("?") | |
| base = f"http://{self.headers.get('Host')}" | |
| if path == "/.well-known/openid-configuration": | |
| self._json( | |
| { | |
| "issuer": base, | |
| "authorization_endpoint": f"{base}/oauth/authorize", | |
| "token_endpoint": f"{base}/oauth/token", | |
| "userinfo_endpoint": f"{base}/oauth/userinfo", | |
| } | |
| ) | |
| return | |
| if path == "/oauth/authorize": | |
| params = urllib.parse.parse_qs(query) | |
| Provider.seen["challenge"] = params["code_challenge"][0] | |
| Provider.seen["nonce"] = params["nonce"][0] | |
| target = params["redirect_uri"][0] + "?" + urllib.parse.urlencode( | |
| {"code": CODE, "state": params["state"][0]} | |
| ) | |
| self.send_response(302) | |
| self.send_header("Location", target) | |
| self.send_header("Content-Length", "0") | |
| self.end_headers() | |
| return | |
| if path == "/oauth/userinfo": | |
| if self.headers.get("Authorization") != f"Bearer {TOKEN}": | |
| self._json({"error": "invalid_token"}, 401) | |
| return | |
| self._json(PEOPLE[Provider.seen["who"]]) | |
| return | |
| self._json({"error": "no such route"}, 404) | |
| def do_POST(self) -> None: # noqa: N802 - framework name | |
| if self.path.split("?", 1)[0] != "/oauth/token": | |
| self._json({"error": "no such route"}, 404) | |
| return | |
| length = int(self.headers.get("Content-Length") or 0) | |
| form = urllib.parse.parse_qs(self.rfile.read(length).decode()) | |
| header = self.headers.get("Authorization") or "" | |
| if not header.startswith("Basic "): | |
| self._json({"error": "invalid_client"}, 401) | |
| return | |
| if base64.b64decode(header.split(" ", 1)[1]).decode() != f"{CLIENT_ID}:{CLIENT_SECRET}": | |
| self._json({"error": "invalid_client"}, 401) | |
| return | |
| verifier = (form.get("code_verifier") or [""])[0] | |
| expected = ( | |
| base64.urlsafe_b64encode(hashlib.sha256(verifier.encode()).digest()) | |
| .decode() | |
| .rstrip("=") | |
| ) | |
| if expected != Provider.seen.get("challenge"): | |
| self._json({"error": "invalid_grant"}, 400) | |
| return | |
| self._json( | |
| { | |
| "access_token": TOKEN, | |
| "token_type": "Bearer", | |
| "id_token": _id_token( | |
| {"sub": PEOPLE[Provider.seen["who"]]["sub"], "nonce": Provider.seen.get("nonce")} | |
| ), | |
| } | |
| ) | |
| def start_provider() -> tuple[str, ThreadingHTTPServer]: | |
| server = ThreadingHTTPServer(("127.0.0.1", 0), Provider) | |
| threading.Thread(target=server.serve_forever, daemon=True).start() | |
| host, port = server.server_address[:2] | |
| return f"http://{host}:{port}", server | |
| def free_port() -> int: | |
| listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM) | |
| listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) | |
| listener.bind(("127.0.0.1", 0)) | |
| port = listener.getsockname()[1] | |
| listener.close() | |
| return port | |
| def wait_for(url: str, *, timeout: float = 180.0) -> None: | |
| import urllib.error | |
| import urllib.request | |
| deadline = time.monotonic() + timeout | |
| while time.monotonic() < deadline: | |
| try: | |
| with urllib.request.urlopen(url, timeout=5): | |
| return | |
| except urllib.error.HTTPError: | |
| return | |
| except (urllib.error.URLError, OSError): | |
| time.sleep(0.4) | |
| raise RuntimeError(f"nothing answered at {url} within {timeout:.0f}s") | |
| class Shots: | |
| """Numbered, captioned screenshots, so the folder reads in order.""" | |
| def __init__(self, directory: Path) -> None: | |
| self.directory = directory | |
| self.count = 0 | |
| self.index: list[tuple[str, str]] = [] | |
| def take(self, page, slug: str, caption: str, *, full: bool = False) -> None: | |
| self.count += 1 | |
| name = f"{self.count:02d}-{slug}.png" | |
| page.screenshot(path=str(self.directory / name), full_page=full) | |
| self.index.append((name, caption)) | |
| print(f" {name} {caption}") | |
| def write_index(self) -> None: | |
| lines = [ | |
| "# Screenshots", | |
| "", | |
| "Taken by `scripts/live_e2e.py`, which starts the real server, a real", | |
| "worker and a real browser and drives them. Two things are stood in for", | |
| "and both are named in that script's docstring: the identity provider,", | |
| "because Hugging Face will not answer an OAuth app that does not exist,", | |
| "and the model, because `--demo-runner` answers without gigabytes of", | |
| "weights. Everything else is the product.", | |
| "", | |
| ] | |
| for name, caption in self.index: | |
| lines.append(f"### {name}") | |
| lines.append("") | |
| lines.append(caption) | |
| lines.append("") | |
| lines.append(f"") | |
| lines.append("") | |
| (self.directory / "README.md").write_text("\n".join(lines), encoding="utf-8") | |
| def sign_in(page, server_url: str, who: str) -> None: | |
| """The real flow: the button, the provider, the callback, the cookie.""" | |
| Provider.seen["who"] = who | |
| page.goto(f"{server_url}/auth/login", wait_until="load") | |
| page.wait_for_timeout(1500) | |
| def enter_the_app(page) -> None: | |
| """Past the splash. The interface opens on a landing view by design.""" | |
| try: | |
| page.get_by_role("button", name=re.compile("Enter the network", re.I)).click(timeout=8000) | |
| page.wait_for_timeout(2500) | |
| except Exception: | |
| pass | |
| def capture(playwright, server_url: str, shots: Shots, worker_started) -> None: | |
| browser = playwright.chromium.launch() | |
| try: | |
| alice = browser.new_context(viewport={"width": 1440, "height": 900}) | |
| page = alice.new_page() | |
| page.goto(server_url, wait_until="load") | |
| page.wait_for_timeout(2000) | |
| shots.take(page, "landing", "The landing page, signed out.") | |
| sign_in(page, server_url, "alice") | |
| enter_the_app(page) | |
| shots.take(page, "alice-signed-in", "Alice, signed in through the real OAuth flow.") | |
| pairing = generate_pairing_code(page) | |
| print(f" pairing code: {pairing}") | |
| shots.take(page, "pairing-code", "Alice generates a one-use pairing code for her machine.") | |
| access_code = worker_started(pairing) | |
| print(f" access code: {access_code}") | |
| page.wait_for_timeout(6000) | |
| page.reload(wait_until="load") | |
| enter_the_app(page) | |
| page.wait_for_timeout(3000) | |
| shots.take(page, "alice-sees-her-worker", "Alice's own worker, visible to her because she paired it.") | |
| bob = browser.new_context(viewport={"width": 1440, "height": 900}) | |
| other = bob.new_page() | |
| sign_in(other, server_url, "bob") | |
| enter_the_app(other) | |
| other.wait_for_timeout(2500) | |
| shots.take( | |
| other, | |
| "bob-signed-in-sees-nothing", | |
| "Bob is signed in and sees no workers at all. Signing in is necessary and not sufficient.", | |
| ) | |
| redeem(other, access_code) | |
| other.wait_for_timeout(4000) | |
| shots.take( | |
| other, | |
| "bob-redeems-the-code", | |
| "Bob pastes the access code the worker printed, and exactly one worker appears.", | |
| ) | |
| submit_a_run(other) | |
| other.wait_for_timeout(12000) | |
| shots.take(other, "run-in-flight", "A request accepted by the worker and running.", full=True) | |
| other.wait_for_timeout(15000) | |
| shots.take(other, "run-complete", "The answer, with the energy this machine reported.", full=True) | |
| open_section(other, "Run a community agent") | |
| other.wait_for_timeout(2000) | |
| # The panel is now three steps rather than a build to download, so the | |
| # screenshot is checked for what it must say. A caption promising a | |
| # Windows build outlived the button that offered one; this asserts | |
| # instead of describing. | |
| setup = other.locator(".c-join").inner_text() | |
| for needle in ("Get the code", "Install it", "pip install -e", "Point it at this server"): | |
| if needle not in setup: | |
| raise AssertionError(f"the setup panel no longer says {needle!r}") | |
| shots.take( | |
| other, | |
| "set-up-the-agent", | |
| "How a volunteer gets the worker running: the code, the install, the command.", | |
| ) | |
| worker_screen(browser, shots) | |
| finally: | |
| browser.close() | |
| def worker_screen(browser, shots: Shots) -> None: | |
| """The worker's own screen, drawn for real and photographed in a browser. | |
| Textual renders to a terminal, and there is no terminal here. Its own SVG | |
| export is the honest way to capture one: it is the same widget tree, laid | |
| out at a real size, with the real text in it. Loading that SVG in the | |
| browser that is already open is simply how it becomes a PNG. | |
| """ | |
| import asyncio | |
| import dataclasses | |
| import threading as _threading | |
| from distinct_agent.dashboard import ActivityLog, JobRow, WorkerView | |
| from distinct_agent.tui import DashboardApp | |
| view = WorkerView( | |
| name="alice-workstation", | |
| agent_id="agent_2a161b0f", | |
| servers=("http://127.0.0.1:7860",), | |
| connected=True, | |
| platform_name="Linux 6.18 (x86_64)", | |
| cpu="cpu", | |
| ram_gb=7.84, | |
| models=("qwen3-0.6b", "olmo-2-1b-instruct"), | |
| tools=(), | |
| queue_capacity=4, | |
| uptime_seconds=47.0, | |
| lifetime_runs=1, | |
| lifetime_measured_runs=1, | |
| lifetime_joules=0.0, | |
| energy_provider="cpu-load-model", | |
| energy_scope="system-cpu-modelled", | |
| containment="none: Linux process restriction is not implemented", | |
| request_sandbox="bubblewrap for the working directory only", | |
| access_code="LRKRX-LFZ2U-CNGL5-4RLAX-ULZQ4-HEJSM-WNWPN-PSTQ", | |
| jobs=(JobRow("job_4f21ab", "running", "qwen3-0.6b", "generating", 0.55, 6.0),), | |
| activity=( | |
| "Paired 'alice-workstation' with http://127.0.0.1:7860", | |
| "Operator approved: qwen3-0.6b, olmo-2-1b-instruct", | |
| "Accepted job_4f21ab", | |
| ), | |
| ) | |
| # Through the log rather than on the view, because that is the path the | |
| # real screen uses: the worker prints, the log catches it, the pane shows | |
| # it. Setting the field directly would photograph a pane that works | |
| # differently from the one that ships. | |
| log = ActivityLog() | |
| for line in view.activity: | |
| log.write(line + "\n") | |
| app = DashboardApp( | |
| worker=None, | |
| base=view, | |
| activity=log, | |
| stop=_threading.Event(), | |
| observer=lambda _w, base, activity=(): dataclasses.replace(base, activity=tuple(activity)), | |
| ) | |
| holder: dict[str, str] = {} | |
| async def render() -> None: | |
| async with app.run_test(size=(120, 34)) as pilot: | |
| await pilot.pause() | |
| await pilot.pause() | |
| holder["svg"] = app.export_screenshot() | |
| await pilot.press("q") | |
| # In its own thread with its own loop. Playwright's synchronous API is | |
| # already running one in this thread, and Textual wants to own the loop it | |
| # runs on; asking either to share is how this deadlocks. | |
| def run_in_a_fresh_loop() -> None: | |
| asyncio.set_event_loop(asyncio.new_event_loop()) | |
| asyncio.get_event_loop().run_until_complete(render()) | |
| thread = _threading.Thread(target=run_in_a_fresh_loop) | |
| thread.start() | |
| thread.join(timeout=120) | |
| if "svg" not in holder: | |
| raise RuntimeError("the worker screen did not render") | |
| page = browser.new_context(viewport={"width": 1200, "height": 760}).new_page() | |
| page.set_content( | |
| '<body style="margin:0;background:#0d1117">' | |
| f'<div style="width:1200px">{holder["svg"]}</div></body>' | |
| ) | |
| page.wait_for_timeout(1200) | |
| shots.take(page, "worker-screen", "The worker's own screen, on the volunteer's machine.") | |
| def open_section(page, title: str) -> None: | |
| """The pairing controls live behind an accordion, which has to be opened.""" | |
| # By role, not by text. The section headings are buttons, and matching on | |
| # text alone finds whichever element happens to contain the words first, | |
| # which was a paragraph. | |
| page.get_by_role("button", name=re.compile(title, re.I)).click(timeout=15000) | |
| page.wait_for_timeout(2500) | |
| def generate_pairing_code(page) -> str: | |
| """Open the section and read the code it mints. | |
| There is no button to press any more: opening the panel is the request, | |
| because attaching a worker is the only reason anybody opens it. This | |
| helper is left with the same name and contract so every caller reads the | |
| same way. | |
| """ | |
| open_section(page, "Run a community agent") | |
| page.wait_for_timeout(3000) | |
| text = page.inner_text("body") | |
| # The code is groups of upper-case characters joined by hyphens. Anchored | |
| # on the length rather than on surrounding words, because the surrounding | |
| # words are the part most likely to be reworded. | |
| found = re.search(r"\b([A-Z0-9]{4,6}(?:-[A-Z0-9]{4,6}){2,})\b", text) | |
| if not found: | |
| raise RuntimeError(f"no pairing code appeared. Page said:\n{text[:1500]}") | |
| return found.group(1) | |
| def redeem(page, access_code: str) -> None: | |
| page.get_by_label(re.compile("Add a worker", re.I)).fill(access_code, timeout=20000) | |
| page.get_by_role("button", name=re.compile("^Add worker$", re.I)).click(timeout=20000) | |
| def submit_a_run(page) -> None: | |
| """Type a request, agree to what it costs, and queue it. | |
| The acknowledgement is not decoration and is not skipped here: the button | |
| is disabled until it is ticked, because the run goes to somebody else's | |
| computer in plain text and the interface says so before it happens. | |
| """ | |
| page.get_by_label(re.compile("Your request", re.I)).fill( | |
| "In one sentence, what is a token?", timeout=20000 | |
| ) | |
| page.wait_for_timeout(500) | |
| page.get_by_role( | |
| "checkbox", name=re.compile("plaintext to a community-operated worker", re.I) | |
| ).check(timeout=20000) | |
| page.wait_for_timeout(500) | |
| page.get_by_role("button", name=re.compile("^Queue run$", re.I)).click(timeout=20000) | |
| def main() -> int: | |
| parser = argparse.ArgumentParser(description="Drive the whole product and photograph it") | |
| parser.add_argument("--shots", type=Path, default=REPO_ROOT / "screenshots") | |
| arguments = parser.parse_args() | |
| from playwright.sync_api import sync_playwright | |
| directory = arguments.shots | |
| if directory.exists(): | |
| shutil.rmtree(directory) | |
| directory.mkdir(parents=True) | |
| shots = Shots(directory) | |
| provider_url, provider = start_provider() | |
| port = free_port() | |
| server_url = f"http://127.0.0.1:{port}" | |
| environment = { | |
| **os.environ, | |
| "DISTINCT_OAUTH_CLIENT_ID": CLIENT_ID, | |
| "DISTINCT_OAUTH_CLIENT_SECRET": CLIENT_SECRET, | |
| "DISTINCT_OAUTH_REDIRECT_URI": f"{server_url}/auth/callback", | |
| "DISTINCT_OAUTH_PROVIDER_URL": provider_url, | |
| "DISTINCT_SESSION_SECRET": "a-fixed-secret-so-this-run-is-repeatable", | |
| "DISTINCT_BIND_HOST": "127.0.0.1", | |
| "PORT": str(port), | |
| "PYTHONPATH": str(REPO_ROOT), | |
| } | |
| environment.pop("DISTINCT_DEV_AUTH", None) | |
| print(f"provider {provider_url}") | |
| print(f"server {server_url}") | |
| server = subprocess.Popen( | |
| [sys.executable, str(REPO_ROOT / "app.py")], | |
| cwd=str(REPO_ROOT), env=environment, | |
| stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, | |
| ) | |
| workers: list[subprocess.Popen] = [] | |
| def start_worker(pairing: str) -> str: | |
| process = subprocess.Popen( | |
| [ | |
| sys.executable, "-m", "distinct_agent", | |
| "--server", server_url, "--pair", pairing, | |
| "--name", "alice-workstation", "--approve", "--demo-runner", | |
| ], | |
| cwd=str(REPO_ROOT), env=environment, | |
| stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, | |
| ) | |
| workers.append(process) | |
| deadline = time.monotonic() + 120 | |
| while time.monotonic() < deadline: | |
| line = process.stdout.readline() | |
| if not line: | |
| if process.poll() is not None: | |
| raise RuntimeError("the worker exited before printing an access code") | |
| continue | |
| print(f" worker | {line.rstrip()}") | |
| found = re.search(r"\b([A-Z0-9]{5}(?:-[A-Z0-9]{4,5}){5,8})\b", line) | |
| if found: | |
| threading.Thread( | |
| target=lambda: [print(f" worker | {line.rstrip()}") for line in process.stdout], | |
| daemon=True, | |
| ).start() | |
| return found.group(1) | |
| raise RuntimeError("the worker never printed an access code") | |
| try: | |
| wait_for(server_url) | |
| print("server is answering") | |
| with sync_playwright() as playwright: | |
| capture(playwright, server_url, shots, start_worker) | |
| shots.write_index() | |
| print(f"\n{shots.count} screenshots in {directory}") | |
| return 0 | |
| finally: | |
| for process in workers: | |
| if process.poll() is None: | |
| process.terminate() | |
| try: | |
| process.wait(timeout=10) | |
| except subprocess.TimeoutExpired: | |
| process.kill() | |
| server.terminate() | |
| try: | |
| server.wait(timeout=15) | |
| except subprocess.TimeoutExpired: | |
| server.kill() | |
| provider.shutdown() | |
| if __name__ == "__main__": | |
| raise SystemExit(main()) | |