diff --git a/ATPROTO.md b/ATPROTO.md index b887e84..4ec523d 100644 --- a/ATPROTO.md +++ b/ATPROTO.md @@ -227,10 +227,12 @@ chart.song.uri -> com.derakkuma.song Publish lexicons: ```bash -ATP_AUTH_TOKEN=... REPO_DID= PDS_URL= \ - deploy/happyview/scripts/publish_lexicon_records.sh +ATP_APP_PASSWORD=... REPO_DID= \ + deploy/happyview/scripts/publish_lexicon_records.py ``` +`publish_lexicon_records.py` uses the same auth conventions as the catalog publisher: it resolves `PDS_URL` from `REPO_DID` when omitted, creates a session from `ATP_APP_PASSWORD`, and still accepts `ATP_AUTH_TOKEN` as an escape hatch. + Publish the catalog: ```bash diff --git a/deploy/happyview/README.md b/deploy/happyview/README.md index 54cb894..48b68a1 100644 --- a/deploy/happyview/README.md +++ b/deploy/happyview/README.md @@ -55,12 +55,12 @@ journalctl -u happyview -f ## Publish lexicon schema records -After obtaining an access token for the lexicon authority account, run from the -Derakkuma repo root: +Run from the Derakkuma repo root with the lexicon authority account's app password: ```bash -ATP_AUTH_TOKEN=... REPO_DID= ./scripts/publish_lexicon_records.sh +ATP_APP_PASSWORD=... REPO_DID= deploy/happyview/scripts/publish_lexicon_records.py ``` This publishes `lexicons/schema-records/com.derakkuma.*.json` as -`com.atproto.lexicon.schema` records. +`com.atproto.lexicon.schema` records. `PDS_URL` is resolved from `REPO_DID` when +omitted; `ATP_AUTH_TOKEN` is still accepted if you already have an access JWT. diff --git a/deploy/happyview/scripts/publish_lexicon_records.py b/deploy/happyview/scripts/publish_lexicon_records.py new file mode 100755 index 0000000..76288ce --- /dev/null +++ b/deploy/happyview/scripts/publish_lexicon_records.py @@ -0,0 +1,63 @@ +#!/usr/bin/env python3 +"""Publish Derakkuma lexicon schema records to an ATProto PDS. + +Environment: + ATP_APP_PASSWORD app password for the lexicon authority account + ATP_IDENTIFIER lexicon authority handle or DID; defaults to REPO_DID + ATP_AUTH_PDS_URL PDS URL for createSession; defaults to PDS_URL / resolved repo PDS + ATP_AUTH_TOKEN optional pre-created access token; skips createSession when set + REPO_DID repo DID to publish lexicon records to (required) + PDS_URL PDS URL for putRecord; if omitted, resolved from REPO_DID + LEXICON_DIR schema record directory; defaults to lexicons/schema-records + ATP_REQUEST_TIMEOUT per-request timeout in seconds; defaults to 30 +""" +from __future__ import annotations + +from pathlib import Path +import json +import os +import sys + +ROOT = Path(__file__).resolve().parents[3] +sys.path.insert(0, str(ROOT / "scripts")) + +from lib.atproto_publish import get_or_create_auth_token, put_record, resolve_pds + +LEXICON_COLLECTION = "com.atproto.lexicon.schema" +DEFAULT_LEXICON_DIR = "lexicons/schema-records" + + +def lexicon_files() -> list[Path]: + lexicon_dir = Path(os.environ.get("LEXICON_DIR", DEFAULT_LEXICON_DIR)) + return sorted(lexicon_dir.glob("com.derakkuma.*.json")) + + +def main() -> int: + repo = os.environ.get("REPO_DID") + if not repo: + print("REPO_DID is required", file=sys.stderr) + return 2 + + pds = os.environ.get("PDS_URL", "").rstrip("/") or resolve_pds(repo) + token = get_or_create_auth_token(repo, pds) + if not token: + return 2 + + files = lexicon_files() + if not files: + print("No lexicon schema record files found", file=sys.stderr) + return 2 + + print(f"Publishing {len(files)} lexicons to {repo} via {pds}") + for index, path in enumerate(files, start=1): + nsid = path.stem + record = json.loads(path.read_text(encoding="utf-8")) + put_record(pds, token, repo, LEXICON_COLLECTION, nsid, record) + print(f"{index}/{len(files)} lexicons published... {nsid}", flush=True) + + print("Lexicon publish complete") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/deploy/happyview/scripts/publish_lexicon_records.sh b/deploy/happyview/scripts/publish_lexicon_records.sh deleted file mode 100755 index 0626236..0000000 --- a/deploy/happyview/scripts/publish_lexicon_records.sh +++ /dev/null @@ -1,23 +0,0 @@ -#!/usr/bin/env bash -set -euo pipefail - -: "${ATP_AUTH_TOKEN:?set ATP_AUTH_TOKEN to an access token for the lexicon authority account}" -: "${REPO_DID:?set REPO_DID to the lexicon authority DID}" -PDS_URL="${PDS_URL:-https://npmx.social}" - -for file in lexicons/schema-records/com.derakkuma.*.json; do - nsid="$(basename "$file" .json)" - body="$(jq -n \ - --arg repo "$REPO_DID" \ - --arg collection "com.atproto.lexicon.schema" \ - --arg rkey "$nsid" \ - --slurpfile record "$file" \ - '{repo:$repo, collection:$collection, rkey:$rkey, record:$record[0]}')" - - echo "publishing $nsid" - curl -fsS -X POST "$PDS_URL/xrpc/com.atproto.repo.putRecord" \ - -H "Authorization: Bearer $ATP_AUTH_TOKEN" \ - -H 'Content-Type: application/json' \ - --data "$body" >/dev/null - echo " ok" -done diff --git a/scripts/lib/__init__.py b/scripts/lib/__init__.py new file mode 100644 index 0000000..8b13789 --- /dev/null +++ b/scripts/lib/__init__.py @@ -0,0 +1 @@ + diff --git a/scripts/lib/atproto_publish.py b/scripts/lib/atproto_publish.py new file mode 100644 index 0000000..d4e832a --- /dev/null +++ b/scripts/lib/atproto_publish.py @@ -0,0 +1,164 @@ +from __future__ import annotations + +import json +import os +import socket +import sys +import time +import urllib.error +import urllib.parse +import urllib.request + + +def request_timeout(env_name: str = "ATP_REQUEST_TIMEOUT", default: str = "30") -> float: + return float(os.environ.get(env_name, os.environ.get("REQUEST_TIMEOUT", default))) + + +def env_int(name: str, default: int) -> int: + return max(1, int(os.environ.get(name, str(default)))) + + +def error_summary(error: BaseException) -> str: + if isinstance(error, urllib.error.HTTPError): + body = "" + try: + body = error.read().decode("utf-8", errors="replace")[:500] + except Exception: # best-effort diagnostics + body = "" + suffix = f" body={body}" if body else "" + return f"HTTP {error.code} {error.reason}{suffix}" + if isinstance(error, urllib.error.URLError): + return f"URL error: {error.reason}" + return f"{error.__class__.__name__}: {error}" + + +def with_retry(operation, attempts: int = 5, label: str = "request"): + for attempt in range(1, attempts + 1): + try: + return operation() + except urllib.error.HTTPError as error: + retryable = error.code == 429 or 500 <= error.code <= 599 + print(f"[{label}] attempt {attempt}/{attempts} failed: {error_summary(error)}", file=sys.stderr, flush=True) + if not retryable or attempt == attempts: + raise + except urllib.error.URLError as error: + print(f"[{label}] attempt {attempt}/{attempts} failed: {error_summary(error)}", file=sys.stderr, flush=True) + if attempt == attempts: + raise + except (TimeoutError, socket.timeout) as error: + print(f"[{label}] attempt {attempt}/{attempts} failed: {error_summary(error)}", file=sys.stderr, flush=True) + if attempt == attempts: + raise + time.sleep(min(2 ** (attempt - 1), 16)) + raise RuntimeError("retry loop exhausted") + + +def get_json(url: str, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> dict: + def request(): + with urllib.request.urlopen(url, timeout=request_timeout(timeout_env)) as resp: # nosec - intentional CLI fetch + return json.loads(resp.read().decode("utf-8")) + + return with_retry(request, label=f"GET {url}") + + +def post_json(url: str, payload: dict, token: str | None = None, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> dict: + body = json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8") + headers = {"Content-Type": "application/json"} + if token: + headers["Authorization"] = f"Bearer {token}" + req = urllib.request.Request(url, data=body, headers=headers, method="POST") + + def request(): + with urllib.request.urlopen(req, timeout=request_timeout(timeout_env)) as resp: # nosec - intentional CLI publish + return json.loads(resp.read().decode("utf-8")) + + return with_retry(request, label=f"POST {url}") + + +def post_bytes(url: str, body: bytes, token: str, mime_type: str, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> dict: + req = urllib.request.Request(url, data=body, headers={"Authorization": f"Bearer {token}", "Content-Type": mime_type}, method="POST") + + def request(): + with urllib.request.urlopen(req, timeout=request_timeout(timeout_env)) as resp: # nosec - intentional CLI publish + return json.loads(resp.read().decode("utf-8")) + + return with_retry(request, label=f"POST {url}") + + +def get_bytes(url: str, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> tuple[bytes, str | None]: + def request(): + with urllib.request.urlopen(url, timeout=request_timeout(timeout_env)) as resp: # nosec - intentional CLI fetch + content_type = resp.headers.get_content_type() + return resp.read(), content_type + + return with_retry(request, label=f"GET {url}") + + +def resolve_pds(did: str, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> str: + doc = get_json(f"https://plc.directory/{urllib.parse.quote(did, safe='')}", timeout_env=timeout_env) + for service in doc.get("service", []): + if service.get("type") == "AtprotoPersonalDataServer": + return str(service["serviceEndpoint"]).rstrip("/") + raise RuntimeError(f"Could not resolve PDS for {did}") + + +def auth_token(auth_pds: str, identifier: str, app_password: str, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> str: + response = post_json( + f"{auth_pds}/xrpc/com.atproto.server.createSession", + {"identifier": identifier, "password": app_password}, + timeout_env=timeout_env, + ) + token = response.get("accessJwt") + if not token: + raise RuntimeError("createSession response did not include accessJwt") + return str(token) + + +def get_or_create_auth_token(repo: str, pds: str, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> str | None: + token = os.environ.get("ATP_AUTH_TOKEN") + if token: + return token + identifier = os.environ.get("ATP_IDENTIFIER", repo) + app_password = os.environ.get("ATP_APP_PASSWORD") or os.environ.get("ATP_PASSWORD") + if not app_password: + print("ATP_APP_PASSWORD is required when ATP_AUTH_TOKEN is not set", file=sys.stderr) + return None + auth_pds = os.environ.get("ATP_AUTH_PDS_URL", "").rstrip("/") or pds + print(f"Authenticating {identifier} via {auth_pds}") + return auth_token(auth_pds, identifier, app_password, timeout_env=timeout_env) + + +def list_records(pds: str, repo: str, collection: str, limit: int = 100, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> list[dict]: + records = [] + cursor = None + while True: + params = {"repo": repo, "collection": collection, "limit": str(limit)} + if cursor: + params["cursor"] = cursor + url = f"{pds}/xrpc/com.atproto.repo.listRecords?{urllib.parse.urlencode(params)}" + page = get_json(url, timeout_env=timeout_env) + records.extend(page.get("records", [])) + cursor = page.get("cursor") + if not cursor: + return records + + +def compact(obj): + if isinstance(obj, dict): + return {k: compact(v) for k, v in obj.items() if v is not None} + if isinstance(obj, list): + return [compact(v) for v in obj] + return obj + + +def put_record(pds: str, token: str, repo: str, collection: str, key: str, record: dict, timeout_env: str = "ATP_REQUEST_TIMEOUT") -> dict: + return post_json( + f"{pds}/xrpc/com.atproto.repo.putRecord", + {"repo": repo, "collection": collection, "rkey": key, "record": compact(record)}, + token, + timeout_env=timeout_env, + ) + + +def rkey_from_uri(uri: str) -> str: + return uri.rstrip("/").split("/")[-1] diff --git a/scripts/publish_catalog.py b/scripts/publish_catalog.py index 9fc4c0c..61d4cf8 100755 --- a/scripts/publish_catalog.py +++ b/scripts/publish_catalog.py @@ -19,15 +19,11 @@ from __future__ import annotations from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timezone import hashlib -import json import mimetypes import os -import socket -import sys -import time -import urllib.error import urllib.parse -import urllib.request + +from lib.atproto_publish import env_int, get_bytes, get_json, get_or_create_auth_token, list_records, post_bytes, put_record, resolve_pds, rkey_from_uri CATALOG_DID = "did:plc:4uoc2as443j2fg2f6xfsogqs" DATA_URL = "https://dp4p6x0xfi5o9.cloudfront.net/maimai/data.json" @@ -81,125 +77,6 @@ def string_value(value: object): return str(value) -def request_timeout() -> float: - return float(os.environ.get("CATALOG_REQUEST_TIMEOUT", "30")) - - -def env_int(name: str, default: int) -> int: - return max(1, int(os.environ.get(name, str(default)))) - - -def error_summary(error: BaseException) -> str: - if isinstance(error, urllib.error.HTTPError): - body = "" - try: - body = error.read().decode("utf-8", errors="replace")[:500] - except Exception: # noqa: BLE001 - best-effort diagnostics - body = "" - suffix = f" body={body}" if body else "" - return f"HTTP {error.code} {error.reason}{suffix}" - if isinstance(error, urllib.error.URLError): - return f"URL error: {error.reason}" - return f"{error.__class__.__name__}: {error}" - - -def with_retry(operation, attempts: int = 5, label: str = "request"): - for attempt in range(1, attempts + 1): - try: - return operation() - except urllib.error.HTTPError as error: - retryable = error.code == 429 or 500 <= error.code <= 599 - print(f"\n[{label}] attempt {attempt}/{attempts} failed: {error_summary(error)}", file=sys.stderr, flush=True) - if not retryable or attempt == attempts: - raise - except urllib.error.URLError as error: - print(f"\n[{label}] attempt {attempt}/{attempts} failed: {error_summary(error)}", file=sys.stderr, flush=True) - if attempt == attempts: - raise - except (TimeoutError, socket.timeout) as error: - print(f"\n[{label}] attempt {attempt}/{attempts} failed: {error_summary(error)}", file=sys.stderr, flush=True) - if attempt == attempts: - raise - time.sleep(min(2 ** (attempt - 1), 16)) - raise RuntimeError("retry loop exhausted") - - -def get_json(url: str) -> dict: - def request(): - with urllib.request.urlopen(url, timeout=request_timeout()) as resp: # nosec - intentional CLI fetch - return json.loads(resp.read().decode("utf-8")) - - return with_retry(request, label=f"GET {url}") - - -def post_json(url: str, payload: dict, token: str | None = None) -> dict: - body = json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8") - headers = {"Content-Type": "application/json"} - if token: - headers["Authorization"] = f"Bearer {token}" - req = urllib.request.Request(url, data=body, headers=headers, method="POST") - - def request(): - with urllib.request.urlopen(req, timeout=request_timeout()) as resp: # nosec - intentional CLI publish - return json.loads(resp.read().decode("utf-8")) - - return with_retry(request, label=f"POST {url}") - - -def post_bytes(url: str, body: bytes, token: str, mime_type: str) -> dict: - req = urllib.request.Request(url, data=body, headers={"Authorization": f"Bearer {token}", "Content-Type": mime_type}, method="POST") - - def request(): - with urllib.request.urlopen(req, timeout=request_timeout()) as resp: # nosec - intentional CLI publish - return json.loads(resp.read().decode("utf-8")) - - return with_retry(request, label=f"POST {url}") - - -def get_bytes(url: str) -> tuple[bytes, str | None]: - def request(): - with urllib.request.urlopen(url, timeout=request_timeout()) as resp: # nosec - intentional CLI fetch - content_type = resp.headers.get_content_type() - return resp.read(), content_type - - return with_retry(request, label=f"GET {url}") - - -def list_records(pds: str, repo: str, collection: str, limit: int = 100) -> list[dict]: - records = [] - cursor = None - while True: - params = {"repo": repo, "collection": collection, "limit": str(limit)} - if cursor: - params["cursor"] = cursor - url = f"{pds}/xrpc/com.atproto.repo.listRecords?{urllib.parse.urlencode(params)}" - page = get_json(url) - records.extend(page.get("records", [])) - cursor = page.get("cursor") - if not cursor: - return records - - -def resolve_pds(did: str) -> str: - doc = get_json(f"https://plc.directory/{urllib.parse.quote(did, safe='')}") - for service in doc.get("service", []): - if service.get("type") == "AtprotoPersonalDataServer": - return str(service["serviceEndpoint"]).rstrip("/") - raise RuntimeError(f"Could not resolve PDS for {did}") - - -def compact(obj): - if isinstance(obj, dict): - return {k: compact(v) for k, v in obj.items() if v is not None} - if isinstance(obj, list): - return [compact(v) for v in obj] - return obj - - -def put_record(pds: str, token: str, repo: str, collection: str, key: str, record: dict) -> dict: - return post_json(f"{pds}/xrpc/com.atproto.repo.putRecord", {"repo": repo, "collection": collection, "rkey": key, "record": compact(record)}, token) - - def cover_url(image_name: str) -> str: base_url = os.environ.get("MAIMAI_COVER_BASE_URL", COVER_BASE_URL).rstrip("/") return f"{base_url}/{urllib.parse.quote(image_name, safe='')}" @@ -209,7 +86,7 @@ def upload_cover_art(pds: str, token: str, image_name: str) -> dict: url = cover_url(image_name) body, content_type = get_bytes(url) mime_type = content_type or mimetypes.guess_type(image_name)[0] or "image/png" - response = post_bytes(f"{pds}/xrpc/com.atproto.repo.uploadBlob", body, token, mime_type) + response = post_bytes(f"{pds}/xrpc/com.atproto.repo.uploadBlob", body, token, mime_type, timeout_env="CATALOG_REQUEST_TIMEOUT") blob = response.get("blob") if not blob: raise RuntimeError(f"uploadBlob response for {image_name} did not include blob") @@ -236,10 +113,6 @@ def publish_parallel(items: list, label: str, concurrency: int, worker, verb: st return results -def rkey_from_uri(uri: str) -> str: - return uri.rstrip("/").split("/")[-1] - - def check_existing_songs(pds: str, repo: str, songs: list[dict], song_rkeys: dict) -> tuple[dict, list[dict]]: print("Checking existing songs") existing_by_rkey = {rkey_from_uri(record.get("uri", "")): record for record in list_records(pds, repo, "com.derakkuma.song")} @@ -276,33 +149,15 @@ def check_existing_charts(pds: str, repo: str, chart_tasks: list[tuple], chart_r return missing_charts -def auth_token(auth_pds: str, identifier: str, app_password: str) -> str: - response = post_json( - f"{auth_pds}/xrpc/com.atproto.server.createSession", - {"identifier": identifier, "password": app_password}, - ) - token = response.get("accessJwt") - if not token: - raise RuntimeError("createSession response did not include accessJwt") - return str(token) - - def main() -> int: repo = os.environ.get("CATALOG_REPO_DID", CATALOG_DID) pds = os.environ.get("CATALOG_PDS_URL", "").rstrip("/") or resolve_pds(repo) fallback_concurrency = os.environ.get("CATALOG_CONCURRENCY") song_concurrency = env_int("CATALOG_SONG_CONCURRENCY", int(fallback_concurrency or "1")) - chart_concurrency = env_int("CATALOG_CHART_CONCURRENCY", int(fallback_concurrency or "8")) - token = os.environ.get("ATP_AUTH_TOKEN") + chart_concurrency = env_int("CATALOG_CHART_CONCURRENCY", int(fallback_concurrency or "2")) + token = get_or_create_auth_token(repo, pds, timeout_env="CATALOG_REQUEST_TIMEOUT") if not token: - identifier = os.environ.get("ATP_IDENTIFIER", repo) - app_password = os.environ.get("ATP_APP_PASSWORD") or os.environ.get("ATP_PASSWORD") - if not app_password: - print("ATP_APP_PASSWORD is required when ATP_AUTH_TOKEN is not set", file=sys.stderr) - return 2 - auth_pds = os.environ.get("ATP_AUTH_PDS_URL", "").rstrip("/") or pds - print(f"Authenticating {identifier} via {auth_pds}") - token = auth_token(auth_pds, identifier, app_password) + return 2 data_url = os.environ.get("MAIMAI_DATA_URL", DATA_URL) now = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") songs = get_json(data_url).get("songs", [])