diff --git a/scripts/catalog_fx_update.py b/scripts/catalog_fx_update.py new file mode 100644 index 00000000..7199aa40 --- /dev/null +++ b/scripts/catalog_fx_update.py @@ -0,0 +1,44 @@ +#!/usr/bin/env python3 +"""Refresh src/treg/catalog/fx.yaml from open.er-api.com (free, keyless). + +Run on deploy or cron: `uv run python scripts/catalog_fx_update.py`. Only currencies already +listed in fx.yaml are refreshed — adding a new billing currency = add its key once, rerun. +""" + +from __future__ import annotations + +import json +import urllib.request +from datetime import date +from pathlib import Path + +import yaml + +FX = Path(__file__).resolve().parent.parent / "src" / "treg" / "catalog" / "fx.yaml" + + +def main() -> int: + doc = yaml.safe_load(FX.read_text()) + live = json.load(urllib.request.urlopen("https://open.er-api.com/v6/latest/USD", timeout=30)) + if live.get("result") != "success": + print(f"rate API answered {live.get('result')!r} — fx.yaml left unchanged") + return 1 + usd_to = live["rates"] + for cur in doc["rates_to_usd"]: + if cur == "USD" or cur not in usd_to: + continue + doc["rates_to_usd"][cur] = round(1 / usd_to[cur], 4) + doc["updated"] = date.today().isoformat() + FX.write_text( + FX.read_text().split("rates_to_usd:")[0] # keep the header comment block + + "rates_to_usd:\n" + + "".join(f" {c}: {v}\n" for c, v in doc["rates_to_usd"].items()) + + f"updated: {doc['updated']}\n" + + f"source: {doc['source']}\n" + ) + print(f"fx.yaml refreshed: {doc['rates_to_usd']} ({doc['updated']})") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/catalog_ingest.py b/scripts/catalog_ingest.py new file mode 100644 index 00000000..85f4634f --- /dev/null +++ b/scripts/catalog_ingest.py @@ -0,0 +1,1826 @@ +#!/usr/bin/env python3 +"""Bulk-ingest a provider's FULL endpoint surface into `src/treg/catalog/.extended.yaml`. + + uv run python scripts/catalog_ingest.py tikhub + uv run python scripts/catalog_ingest.py dataforseo justoneapi + uv run python scripts/catalog_ingest.py all --refresh + +The hand-curated `.yaml` (tier: core) stays the quality tier: capability-mapped, live +verified, with captured example responses. This script produces the *coverage* tier — every route +the provider exposes, machine-generated, unmapped and unverified — so an agent browsing the +marketplace sees the real breadth of what a key unlocks. + +Properties this script guarantees (they are why it exists instead of a one-off scrape): + - **re-runnable**: same upstream in, byte-identical file out. Endpoints sort by (platform, path); + nothing carries a timestamp except `source.ingested`, which only moves when the data does. + - **core wins**: any (method, path) already in the provider's core yaml is skipped, never duplicated. + - **cached**: downloads land in ~/.cache/treg-catalog-ingest (override TREG_INGEST_CACHE); + `--refresh` re-fetches. Nothing is cached inside the repo. + - **credential-free**: no source below needs auth, so no secret can reach a committed file. + +See docs/context/architecture/catalog.md for the extended-entry schema. +""" + +from __future__ import annotations + +import argparse +import json +import os +import re +import subprocess +import sys +import time +from concurrent.futures import ThreadPoolExecutor +from datetime import date +from pathlib import Path + +import yaml + +ROOT = Path(__file__).resolve().parent.parent +CATALOG = ROOT / "src" / "treg" / "catalog" +CACHE = Path(os.environ.get("TREG_INGEST_CACHE") or (Path.home() / ".cache" / "treg-catalog-ingest")) + +CJK = re.compile("[ -鿿＀-￯]") + + +class _TolerantLoader(yaml.SafeLoader): + """SafeLoader that survives real-world specs — DataForSEO's has a bare `=` enum value, which + YAML 1.1 reserves as the 'default value' tag and SafeLoader refuses to construct.""" + + +_TolerantLoader.add_constructor("tag:yaml.org,2002:value", lambda loader, node: node.value) + + +# --------------------------------------------------------------------------------------------- +# plumbing + + +def fetch(url: str, key: str, *, refresh: bool = False, headers: dict[str, str] | None = None) -> bytes: + """GET `url` through curl, cached on disk by `key`. Tolerates a truncated body (see tikhub).""" + CACHE.mkdir(parents=True, exist_ok=True) + blob = CACHE / key + if blob.is_file() and not refresh: + return blob.read_bytes() + cmd = ["curl", "-sSL", "--max-time", "300", "-o", str(blob)] + for k, v in (headers or {}).items(): + cmd += ["-H", f"{k}: {v}"] + cmd.append(url) + proc = subprocess.run(cmd, capture_output=True) + if not blob.is_file() or blob.stat().st_size == 0: + raise SystemExit(f"fetch failed: {url}\n{proc.stderr.decode()[:400]}") + return blob.read_bytes() + + +def salvage_json_map(text: str, at_key: str) -> dict: + """Parse the `at_key` object of a JSON document that may be TRUNCATED mid-stream. + + api.tikhub.io/openapi.json is cut off by the origin at ~840 KB (reproducible across encodings + and clients), so a plain json.loads yields nothing at all. Decoding entry-by-entry keeps every + complete member and drops only the one that was cut. + """ + dec = json.JSONDecoder() + start = text.index(f'"{at_key}"') + pos = text.index("{", start) + 1 + out: dict = {} + while True: + while pos < len(text) and text[pos] in " \t\r\n,": + pos += 1 + if pos >= len(text) or text[pos] != '"': + break + try: + key, pos = dec.raw_decode(text, pos) + pos = text.index(":", pos) + 1 + while text[pos] in " \t\r\n": + pos += 1 + val, pos = dec.raw_decode(text, pos) + except (ValueError, IndexError): + break + out[key] = val + return out + + +ZERO_WIDTH = re.compile("[​-‏‪-‮]") + + +def clean(s: str) -> str: + """Collapse whitespace and drop the zero-width joiners providers sprinkle through their docs.""" + return " ".join(ZERO_WIDTH.sub("", s or "").split()).strip() + + +def english(s: str) -> str: + """TikHub summaries are '中文描述/English description' — keep the English half.""" + s = clean(s) + parts = s.split("/") + for i in range(len(parts)): + cand = "/".join(parts[i:]).strip() + if cand and not CJK.search(cand): + return cand + return s + + +def titlecase(slug: str) -> str: + words = re.split(r"[_\-]+", slug) + return " ".join(words).strip().capitalize() + + +def short_name(candidate: str, summary: str) -> str: + """`candidate` as an endpoint's short display `name`, or "" when it adds nothing. + + A name is only worth carrying when the upstream spec offers a real human title that is + DISTINCT from the description we keep verbatim in `summary`: non-empty after cleaning, no + CJK left, fits a row heading (≤60 chars — the validator's cap), and not just the summary + (or its first 60 chars) restated. + """ + cand = english(candidate or "") + if not cand or CJK.search(cand) or len(cand) > 60: + return "" + if cand.lower() in (clean(summary or "").lower(), clean(summary or "")[:60].lower()): + return "" + return cand + + +def slug_id(provider: str, path: str) -> str: + body = re.sub(r"[^a-z0-9]+", "-", path.lower().lstrip("/")).strip("-") + return f"{provider}.x.{body}" + + +def core_routes(provider: str) -> set[tuple[str, str]]: + """(METHOD, path) pairs already curated in the provider's core yaml — those win.""" + f = CATALOG / f"{provider}.yaml" + if not f.is_file(): + return set() + data = yaml.safe_load(f.read_text()) or {} + return { + (str(ep.get("method", "")).upper(), str(ep.get("path", ""))) + for ep in (data.get("endpoints") or []) + } + + +def carry_verification(provider: str, endpoints: list[dict]) -> int: + """Re-attach `verified` / `example_response` / `unverified` / `capability` from the file being replaced. + + Those fields are the only ones NOT derived from upstream: verification stamps are the result + of an actual paid call made by scripts/catalog_verify_extended.py, and `capability` is a + reviewed mapping into the capabilities.yaml taxonomy (2026-07-28 batch onwards). `name` (the + short display title) is carried on the same guard: ingest generates one where the spec offers + it, but a reviewed/hand-set title must survive a re-ingest just like a reviewed capability. + Regenerating the file must not silently discard them, so they are carried across by id, as long as the + route itself (method + path) still matches — a route that moved is a different endpoint and + its old result means nothing. + + A carried `capability` brings its `platform` with it: the validator requires platform == + the capability's first segment, and review sometimes refines the ingest guess (e.g. douyin -> + douyin-xingtu, tiktok -> tiktok-shop), so regenerating the guess over the reviewed platform + would break validation on the next run. + + A stamp carries its `test_request` with it. The generated request is our best guess from the + docs; a request stamped `verified` is one we watched succeed, including the case where the + verifier had to repair it (a page-size knob the endpoint rejected). Regenerating the guess over + the top of a proven request would drop the stamp on the next ingest and re-bill the call to + learn what the file already knew. `--refresh` re-fetches the docs; re-verifying is what + replaces a proven request, not re-ingesting. + """ + old_file = CATALOG / f"{provider}.extended.yaml" + if not old_file.is_file(): + return 0 + old = {ep["id"]: ep for ep in (yaml.safe_load(old_file.read_text()) or {}).get("endpoints") or []} + kept = 0 + for ep in endpoints: + prev = old.get(ep["id"]) + if not prev: + continue + if prev.get("method") != ep.get("method") or prev.get("path") != ep.get("path"): + continue + for field in ("verified", "example_response", "unverified", "capability", "name"): + if prev.get(field) is not None: + ep[field] = prev[field] + if prev.get("capability") is not None and prev.get("platform"): + # the reviewed platform (capability's first segment) wins over the ingest guess + ep["platform"] = prev["platform"] + if prev.get("verified"): + kept += 1 + if prev.get("test_request") is not None: + ep["test_request"] = prev["test_request"] + ep.pop("untestable", None) + return kept + + +def write_extended(provider: str, source: dict, endpoints: list[dict], notes: list[str]) -> Path: + carried = carry_verification(provider, endpoints) + if carried: + print(f" carried {carried} verification stamp(s) forward", file=sys.stderr) + endpoints = sorted(endpoints, key=lambda e: (e["platform"], e["path"], e["method"])) + seen: set[str] = set() + for ep in endpoints: + if ep["id"] in seen: + raise SystemExit(f"{provider}: duplicate generated id {ep['id']}") + seen.add(ep["id"]) + out = CATALOG / f"{provider}.extended.yaml" + header = [ + f"# {provider} — EXTENDED tier: the provider's full endpoint surface, machine-generated by", + "# scripts/catalog_ingest.py. Do not hand-edit; re-run the script instead. Entries start with", + "# no capability, no verification and no example response — verification stamps and reviewed", + "# `capability` mappings (with their platform correction) are added later and carried across", + f"# re-ingests by id via carry_verification. Routes curated in core {provider}.yaml are excluded here.", + ] + header += [f"# {n}" for n in notes] + body = yaml.safe_dump( + {"provider": provider, "source": source, "endpoints": endpoints}, + sort_keys=False, + allow_unicode=True, + width=4096, + default_flow_style=False, + ) + out.write_text("\n".join(header) + "\n\n" + body) + return out + + +# --------------------------------------------------------------------------------------------- +# tikhub + +# path segment after /api/v1/ → platform slug. Absent = use the segment as-is. +TIKHUB_PLATFORM = {"twitter": "x", "net_ease_cloud_music": "netease-music", "hybrid": "web"} +TIKHUB_PLATFORM.update({k: k.replace("_", "-") for k in ("wechat_channels", "wechat_mp", "wechat_search")}) +# utility families, not data: solved captchas, throwaway inboxes, the account's own meter, the +# iOS-shortcut installer, the MCP bridge and the Sora2 video generator. +TIKHUB_SKIP = {"captcha", "temp_mail", "ios_shortcut", "mcp", "sora2", "tikhub", "health", "downloader"} + + +# --- TikHub parameter documentation ----------------------------------------------------------- +# +# api.tikhub.io/openapi.json is served truncated at 842001 bytes (see the note in ingest_tikhub), +# which costs us the parameter block of ~78% of the routes. TikHub's published docs are an Apifox +# site, and Apifox serves the same project as JSON through the API its own front-end calls — no +# auth, no truncation, and richer than the OpenAPI: every parameter carries a `sampleValue`, the +# working example value TikHub themselves demo the endpoint with. That is what makes a generated +# `test_request` a real call rather than a guess. +# +# /api/v1/published-projects/domains/ → the project + branch ids +# /api/v1/published-projects//http-api-tree → every documented operation, with its api id +# /api/v1/published-projects//http-apis/ → one operation: parameters, requestBody, samples + +TIKHUB_DOCS = "https://docs.tikhub.io" +_JSON_HDR = {"accept": "application/json"} + +# Parameters that name the CALLER's own infrastructure or credentials. A test request must never +# invent a value for these — an endpoint that truly needs one cannot be verified with a bare key. +TIKHUB_PARAM_SKIP = { + "cookie", "cookies", "proxy", "proxy_url", "token", "api_key", "apikey", "authorization", + "x-api-key", "session", "sessionid", "auth", "access_token", "refresh_token", +} +# Page-size knobs. Clamped to the smallest useful page: it keeps the captured example small and, +# on the per-result endpoints, keeps the call cheap. +TIKHUB_COUNT_PARAMS = { + "count", "limit", "size", "num", "number", "page_size", "pagesize", "per_page", "perpage", + "page_count", "pagecount", "page_size_", "result_count", +} +TIKHUB_TEST_COUNT = 5 + + +def tikhub_docs_apis(refresh: bool) -> dict[str, dict]: + """`path` → the Apifox operation object (parameters, requestBody, sampleValues).""" + proj = json.loads(fetch( + f"{TIKHUB_DOCS}/api/v1/published-projects/domains/docs.tikhub.io", + "tikhub_apifox_project.json", refresh=refresh, headers=_JSON_HDR, + ))["data"] + pid = proj["id"] + tree = json.loads(fetch( + f"{TIKHUB_DOCS}/api/v1/published-projects/{pid}/http-api-tree?locale=en-US", + "tikhub_apifox_tree.json", refresh=refresh, headers=_JSON_HDR, + ))["data"] + + ids: list[int] = [] + + def walk(nodes: list) -> None: + for n in nodes: + api = n.get("api") or {} + if api.get("id") and api.get("path"): + ids.append(api["id"]) + walk(n.get("children") or []) + + walk(tree) + todo = [i for i in ids if refresh or not (CACHE / "tikhub_apifox" / f"{i}.json").is_file()] + if todo: + print(f" fetching {len(todo)} documented operation(s) from the Apifox docs API…", file=sys.stderr) + + def one(api_id: int) -> None: + fetch(f"{TIKHUB_DOCS}/api/v1/published-projects/{pid}/http-apis/{api_id}?locale=en-US", + f"tikhub_apifox/{api_id}.json", refresh=refresh, headers=_JSON_HDR) + + (CACHE / "tikhub_apifox").mkdir(parents=True, exist_ok=True) + with ThreadPoolExecutor(max_workers=8) as pool: + list(pool.map(one, todo)) + + out: dict[str, dict] = {} + for api_id in ids: + blob = CACHE / "tikhub_apifox" / f"{api_id}.json" + if not blob.is_file(): + continue + try: + data = json.loads(blob.read_text())["data"] + except (ValueError, KeyError): + continue + if data.get("path"): + out[data["path"]] = data + return out + + +def tikhub_docs_schemas(refresh: bool) -> dict[str, dict]: + """`#/definitions/` targets — the shared request/response models a requestBody may $ref.""" + proj = json.loads(fetch( + f"{TIKHUB_DOCS}/api/v1/published-projects/domains/docs.tikhub.io", + "tikhub_apifox_project.json", refresh=refresh, headers=_JSON_HDR, + ))["data"] + raw = json.loads(fetch( + f"{TIKHUB_DOCS}/api/v1/published-projects/{proj['id']}/data-schemas", + "tikhub_apifox_schemas.json", refresh=refresh, headers=_JSON_HDR, + ))["data"] + return {str(s["id"]): (s.get("jsonSchema") or {}) for s in raw if s.get("id")} + + +def _schema_example(schema: dict, defs: dict, depth: int = 0): + """A concrete value for `schema`, built only from values the docs actually state. + + Returns None when the docs state none — a request body we invented would test our imagination, + not the endpoint. Arrays are built with a SINGLE item: these are batch routes, and on a + per-result price a 10-id example body costs 10x for no extra information. + """ + if depth > 4 or not isinstance(schema, dict): + return None + ref = schema.get("$ref") + if ref: + return _schema_example(defs.get(str(ref).rsplit("/", 1)[-1]) or {}, defs, depth + 1) + for key in ("default", "example"): + if schema.get(key) not in (None, "", [], {}): + return schema[key] + examples = schema.get("examples") or [] + if examples and examples[0] not in (None, ""): + return examples[0] + stype = schema.get("type") + if stype == "array": + item = _schema_example(schema.get("items") or {}, defs, depth + 1) + return None if item is None else [item] + if stype == "object": + props = schema.get("properties") or {} + required = set(schema.get("required") or []) + out = {} + for name, sub in props.items(): + val = _schema_example(sub if isinstance(sub, dict) else {}, defs, depth + 1) + if val is not None: + out[name] = val + elif name in required: + return None # a required field we cannot fill makes the whole body a guess + return out + return None + + +def _body_example(body_spec: dict, defs: dict): + """The request body a test call should send, or None if the docs give us nothing usable.""" + for ex in body_spec.get("examples") or []: + value = ex.get("value") + if isinstance(value, str): + try: + value = json.loads(value) + except ValueError: + continue + if value in (None, ""): + continue + return value[:1] if isinstance(value, list) else value + built = _schema_example(body_spec.get("jsonSchema") or {}, defs) + if built is None and not body_spec.get("required"): + return {} + return built + + +def _param_value(p: dict): + """The value a test request should send for parameter `p`, or None if the docs give us none. + + Order: TikHub's own demo value → the schema default → the first enum member. Anything else + would be us inventing an id, and an invented id verifies nothing. + """ + schema = p.get("schema") or {} + for cand in (p.get("sampleValue"), schema.get("default")): + if cand not in (None, "", [], {}): + return cand + enum = schema.get("enum") or [] + if enum: + return enum[0] + return None + + +def _coerce(value, ptype: str): + if ptype == "integer" and isinstance(value, str) and value.lstrip("-").isdigit(): + return int(value) + if ptype == "boolean" and isinstance(value, str): + return value.strip().lower() in ("true", "1", "yes") + return value + + +def tikhub_input_and_test(op: dict, defs: dict) -> tuple[dict, dict, str]: + """Build an endpoint's `input` schema and `test_request` from one Apifox operation. + + Returns (input, test_request, reason). `test_request` is empty when the operation cannot be + called blind — every such case gets a `reason` that ends up in the file as an `unverified:` + note, so a human can see WHY rather than just that it is missing. + + Which parameters the test sends: + - every REQUIRED one (if any of them has no documented value, there is no test request — + guessing a user id produces a 404 that says nothing about the endpoint); + - every optional one that carries a `sampleValue`, because TikHub marks the real selector + optional all the time (fetch_user_profile's `uniqueId` is "optional" but the call is + meaningless without it) — a plain schema default is left off instead, since the server + applies it anyway; + - page-size knobs clamped to TIKHUB_TEST_COUNT. + """ + params = (op.get("parameters") or {}).get("query") or [] + body_spec = op.get("requestBody") or {} + inp: dict = {} + qs: dict = {} + test_q: dict = {} + unresolved: list[str] = [] + skipped: list[str] = [] + + for p in params: + name = p.get("name") + if not name or p.get("enable") is False: + continue + schema = p.get("schema") or {} + ptype = p.get("type") or schema.get("type") or "string" + entry: dict = {"type": ptype, "required": bool(p.get("required"))} + note = english(p.get("description") or "") + if note: + entry["note"] = note + value = _param_value(p) + if value not in (None, "", [], {}): + entry["example"] = _coerce(value, ptype) + if schema.get("enum"): + entry["enum"] = schema["enum"] + qs[name] = entry + + if name.lower() in TIKHUB_PARAM_SKIP: + if p.get("required"): + skipped.append(name) + continue + if name.lower() in TIKHUB_COUNT_PARAMS and ptype in ("integer", "number"): + test_q[name] = TIKHUB_TEST_COUNT + continue + if value in (None, "", [], {}): + if p.get("required"): + unresolved.append(name) + continue + if p.get("required") or p.get("sampleValue") not in (None, "", [], {}): + test_q[name] = _coerce(value, ptype) + + if qs: + inp["queryParams"] = qs + + body = None + if (body_spec.get("type") or "none") == "application/json": + js = body_spec.get("jsonSchema") or {} + if js.get("$ref"): + js = defs.get(str(js["$ref"]).rsplit("/", 1)[-1]) or js + inp["bodyType"] = "json" + inp["body"] = {"type": js.get("type", "object"), "required": bool(body_spec.get("required"))} + note = english(js.get("description") or "") + if note: + inp["body"]["note"] = note + props = js.get("properties") or {} + if props: + inp["body"]["properties"] = { + k: {"type": v.get("type", "string"), + **({"note": english(v["description"])} if isinstance(v, dict) and v.get("description") else {}), + **({"required": True} if k in set(js.get("required") or []) else {})} + for k, v in props.items() if isinstance(v, dict) + } + body = _body_example(body_spec, defs) + if body is None: + unresolved.append("") + # the same rule as for query params: a body field naming the caller's own cookie/token + # cannot be filled by us, and the docs' placeholder ("your_vip_bilibili_cookie") would send + # a call that fails for a reason the failure text would not explain + elif isinstance(body, dict): + own = sorted(k for k in body if k.lower() in TIKHUB_PARAM_SKIP) + if own: + skipped.extend(own) + + if skipped: + return inp, {}, f"needs the caller's own {', '.join(skipped)}" + if unresolved: + return inp, {}, f"no documented value for required {', '.join(unresolved)}" + + test: dict = {} + if test_q: + test["queryParams"] = test_q + if body is not None: + test["body"] = body + if not test: + # no parameters at all: the route is callable bare, and an empty test_request is the + # honest description of that call + test = {"queryParams": {}} + return inp, test, "" + + +def ingest_tikhub(refresh: bool) -> tuple[Path, dict]: + card = json.loads(fetch( + "https://api.tikhub.io/api/v1/tikhub/user/get_all_endpoints_info", + "tikhub_rate_card.json", refresh=refresh, + ))["data"] + spec_text = fetch("https://api.tikhub.io/openapi.json", "tikhub_openapi.json", refresh=refresh) + spec_paths = salvage_json_map(spec_text.decode("utf-8", "replace"), "paths") + + routes: dict[str, dict] = {} + for row in card: + # the rate card carries a few malformed uris — a stray \r\n, one missing its leading slash + uri = row["endpoint_uri"].strip() + if not uri.startswith("/"): + uri = "/" + uri + seg = uri.strip("/").split("/") + if len(seg) < 4 or seg[0] != "api" or seg[1] != "v1" or seg[2] in TIKHUB_SKIP: + continue + routes.setdefault(uri, row) + + methods = tikhub_methods(sorted(routes), refresh=refresh) + docs_apis = tikhub_docs_apis(refresh) + docs_defs = tikhub_docs_schemas(refresh) + skip = core_routes("tikhub") + endpoints, from_spec, with_test, documented, dropped_undocumented = [], 0, 0, 0, 0 + for uri, row in routes.items(): + family = uri.strip("/").split("/")[2] + platform = TIKHUB_PLATFORM.get(family, family.replace("_", "-")) + op = spec_paths.get(uri) or {} + doc_op = docs_apis.get(uri) + method = (methods.get(uri) + or (doc_op or {}).get("method", "").upper() + or next((m.upper() for m in ("get", "post") if m in op), "GET")) + if (method, uri) in skip: + continue + label = titlecase(family) + spec_summary = english((op.get(method.lower()) or {}).get("summary") or "") + if CJK.search(spec_summary): + spec_summary = "" # a handful of routes are documented in Chinese only + if spec_summary: + from_spec += 1 + # TikHub's terse ones ("Home Feed", "Search user") lose their subject once the + # endpoint is listed next to 22 other platforms — give it back + summary = spec_summary if len(spec_summary) >= 16 else f"{label}: {spec_summary}" + else: + tail = uri.strip("/").split("/")[3:] + where = "/".join(tail[:-1]) + summary = f"{label}: {titlecase(tail[-1]).lower()}" + (f" ({where})" if where else "") + cost = float(row.get("endpoint_cost") or 0) + entry = { + "id": slug_id("tikhub", uri.replace("/api/v1/", "")), + "tier": "extended", + "platform": platform, + "method": method, + "path": uri, + "summary": summary, + } + # Apifox gives every documented op a human title ("中文/Get TikHub user info") distinct + # from the openapi summary — the English half is the display `name`. + title = short_name((doc_op or {}).get("name") or "", summary) + if title: + entry["name"] = title + entry["cost"] = ( + {"type": "free", "value": 0.0, "currency": "USD"} if cost == 0 + else {"type": "per_success", "value": cost, "currency": "USD"} + ) + if doc_op is not None: + documented += 1 + inp, test, reason = tikhub_input_and_test(doc_op, docs_defs) + if inp: + entry["input"] = inp + if test: + entry["test_request"] = test + with_test += 1 + elif reason: + entry["untestable"] = reason + else: + # A route the rate card lists but the docs never describe is not usable inventory: + # no parameter spec exists anywhere TikHub publishes, and spot checks 404 ("ghost + # routes" the server enumerates but no longer serves). Dropped rather than catalogued — + # 399 such routes on 2026-07-28; they return automatically if TikHub ever documents them. + dropped_undocumented += 1 + continue + endpoints.append(entry) + + source = { + "method": "openapi + provider rate card", + "ingested": str(date.today()), + "spec_urls": [ + "https://api.tikhub.io/openapi.json", + "https://api.tikhub.io/api/v1/tikhub/user/get_all_endpoints_info", + "https://docs.tikhub.io/api/v1/published-projects/{id}/http-api-tree", + ], + } + notes = [ + "", + "Route list + per-endpoint USD price come from TikHub's own free rate card endpoint", + "(get_all_endpoints_info), which is authoritative and complete. openapi.json is served", + f"TRUNCATED by the origin (~840 KB cap), so it supplied summaries for only {from_spec} routes;", + "the rest have a summary derived from the route itself. HTTP methods are ground truth: an", + "OPTIONS request answers 405 with an `allow:` header without executing (or billing) the route.", + "Skipped families: captcha, temp_mail, ios_shortcut, sora2, tikhub (own account), health.", + "", + f"`input` and `test_request` come from TikHub's Apifox docs API, which documents {documented}", + "of these routes in full — parameter types, which are required, and a `sampleValue` per", + f"parameter that is TikHub's own working demo value. {with_test} routes got a test request;", + "the rest carry `untestable:` saying what is missing (a parameter only the caller can", + f"supply, such as their own platform cookie). {dropped_undocumented} rate-card routes with no", + "documentation anywhere were DROPPED entirely (ghost routes; several 404 live). `verified` +", + "`example_response` are stamped later by scripts/catalog_verify_extended.py and are carried", + "across re-ingests by id, so regenerating this file does not throw away live test results.", + ] + return write_extended("tikhub", source, endpoints, notes), {"from_spec": from_spec} + + +def tikhub_methods(paths: list[str], *, refresh: bool) -> dict[str, str]: + """Ground-truth HTTP method per route, via an OPTIONS probe (405 + `allow:` header). + + Starlette answers a non-matching method with 405 before the handler runs, so the probe is free + — critical here, because a bare GET against a route whose params are all optional would execute + and bill (the Moz quota trap in catalog.md, at 1400x). + + A probe that comes back without an `allow:` header (transient 5xx, connection reset — TikHub + does both under load) is a MISS, not an answer: it is retried in-process and never written to + the cache. Caching misses is how the first run of this script left 1361 of 1404 routes falling + back to the "GET" default, 128 of them wrongly (they are POST-only). + """ + CACHE.mkdir(parents=True, exist_ok=True) + blob = CACHE / "tikhub_methods.json" + known: dict[str, str] = {} + if blob.is_file() and not refresh: + known = {k: v for k, v in json.loads(blob.read_text()).items() if v} + todo = [p for p in paths if p not in known] + if todo: + print(f" probing {len(todo)} route method(s) via OPTIONS…", file=sys.stderr) + + def probe(path: str) -> tuple[str, str]: + for attempt in range(3): + out = subprocess.run( + ["curl", "-sS", "-X", "OPTIONS", "-D", "-", "-o", "/dev/null", + "--max-time", "30", f"https://api.tikhub.io{path}"], + capture_output=True, + ).stdout.decode("utf-8", "replace") + m = re.search(r"^allow:\s*(.+)$", out, re.I | re.M) + allowed = [a.strip().upper() for a in (m.group(1) if m else "").split(",") if a.strip()] + for pref in ("GET", "POST", "PUT", "PATCH", "DELETE"): + if pref in allowed: + return path, pref + time.sleep(1 + attempt) + return path, "" + + with ThreadPoolExecutor(max_workers=8) as pool: + for path, method in pool.map(probe, todo): + if method: + known[path] = method + blob.write_text(json.dumps(known, indent=0, sort_keys=True)) + return {k: v for k, v in known.items() if v} + + +# --------------------------------------------------------------------------------------------- +# dataforseo + +# /v3///… → the platform the data is ABOUT (not the API family it lives under). +DFS_PLATFORM = { + ("serp", "google"): "google", ("serp", "bing"): "bing", ("serp", "yahoo"): "yahoo", + ("serp", "baidu"): "baidu", ("serp", "naver"): "naver", ("serp", "seznam"): "seznam", + ("serp", "youtube"): "youtube", ("serp", "ai_summary"): "google", ("serp", "screenshot"): "google", + ("keywords_data", "google_ads"): "google", ("keywords_data", "google_trends"): "google", + ("keywords_data", "bing"): "bing", + ("dataforseo_labs", "google"): "google", ("dataforseo_labs", "bing"): "bing", + ("dataforseo_labs", "amazon"): "amazon", ("dataforseo_labs", "apple"): "app-store", + ("business_data", "google"): "google", ("business_data", "tripadvisor"): "tripadvisor", + ("business_data", "trustpilot"): "trustpilot", + ("merchant", "amazon"): "amazon", ("merchant", "google"): "google", + ("app_data", "google"): "google-play", ("app_data", "apple"): "app-store", + ("ai_optimization", "chat_gpt"): "chatgpt", ("ai_optimization", "gemini"): "gemini", + ("ai_optimization", "claude"): "claude", ("ai_optimization", "perplexity"): "perplexity", + ("ai_optimization", "llm_mentions"): "ai-search", + ("ai_optimization", "ai_keyword_data"): "ai-search", +} +# result readers and task bookkeeping — not a data call an agent would choose +DFS_PLUMBING = { + "tasks_ready", "tasks_fixed", "id_list", "errors", "webhook_resend", "status", + "available_filters", "locations", "languages", "locations_and_languages", "categories", + "models", "versions", "user_data", "force_stop", "task_get", +} +DFS_DOCS = re.compile(r"https://docs\.dataforseo\.com/[^'\s\"]+") + + +def ingest_dataforseo(refresh: bool) -> tuple[Path, dict]: + raw = fetch( + "https://raw.githubusercontent.com/dataforseo/OpenApiDocumentation/master/openapi_specification.yaml", + "dataforseo_openapi.yaml", refresh=refresh, + ) + spec = yaml.load(raw, Loader=_TolerantLoader) + ops: dict[str, tuple[str, dict]] = {} + for path, item in spec["paths"].items(): + for method, op in item.items(): + if method.lower() in ("get", "post", "put", "patch", "delete"): + ops[path] = (method.upper(), op) + break + + # one canonical call per operation: the sync `live` form when it exists, else `task_post`. + families: dict[str, list[str]] = {} + for path in ops: + seg = path.strip("/").split("/") + if seg[-1] in DFS_PLUMBING or "task_get" in seg or path.endswith("/html"): + continue + stem = "/".join(seg[: seg.index("live")]) if "live" in seg else ( + "/".join(seg[:-1]) if seg[-1] == "task_post" else path) + families.setdefault(stem, []).append(path) + chosen: list[str] = [] + for paths in families.values(): + live = [p for p in paths if "/live" in p] + chosen += live or paths + + skip = core_routes("dataforseo") + endpoints = [] + for path in chosen: + method, op = ops[path] + if (method, path) in skip: + continue + seg = path.strip("/").split("/") + platform = DFS_PLATFORM.get((seg[1], seg[2] if len(seg) > 2 else ""), "web") + desc = clean(op.get("description") or "") + docs = DFS_DOCS.search(desc) + desc = re.split(r"\s*for more info please visit", desc)[0] + summary = (re.split(r"(?<=[.!?])\s+", desc)[0] or op.get("operationId") or seg[-1]).strip() + if not summary: + summary = titlecase(seg[-1]) + entry = { + "id": slug_id("dataforseo", path.replace("/v3/", "")), + "tier": "extended", + "platform": platform, + "method": method, + "path": path, + "summary": summary[:400], + } + # DataForSEO's spec has NO per-op summary/title field (checked 2026-07-28: 0 of 570 ops) + # — the only short handle is the CamelCase operationId ("GoogleOrganicLiveRegular"), + # which de-camels into a serviceable display title. + title = short_name(re.sub(r"(?<=[a-z0-9])(?=[A-Z])", " ", op.get("operationId") or ""), summary) + if title: + entry["name"] = title + if docs: + entry["docs_url"] = docs.group(0).rstrip("',.").split("?")[0] + endpoints.append(entry) + + source = { + "method": "openapi", + "ingested": str(date.today()), + "spec_urls": [ + "https://github.com/dataforseo/OpenApiDocumentation/blob/master/openapi_specification.yaml" + ], + } + notes = [ + "", + "Selection rule: one canonical call per operation. DataForSEO exposes most jobs three ways", + "(sync `/live/...`, async `/task_post` + `/task_get`, plus a `/live/html` raw-HTML twin);", + "the sync `live` form is kept when it exists, `task_post` only when it does not. Task-", + "bookkeeping routes (tasks_ready, tasks_fixed, id_list, errors, task_get) and reference", + "lookups (locations, languages, categories, available_filters) are omitted.", + "`platform` is the system the data DESCRIBES (google, amazon, app-store, chatgpt…), not the", + "API family; anything not tied to one engine is `web`. No per-endpoint price: DataForSEO's", + "pricing is published per API family on a separate page, not in the spec.", + ] + return write_extended("dataforseo", source, endpoints, notes), {} + + +# --------------------------------------------------------------------------------------------- +# justoneapi + +# Just One API's own x-platform-id values are marketing names ("douyin-tiktok-china"); normalise +# them onto the slugs the rest of the catalog already uses so a platform page shows every provider. +JOA_PLATFORM = { + "douban-movie": "douban", + "dewu-poizon": "dewu", + "douyin-tiktok-china": "douyin", + "douyin-e-commerce": "douyin-shop", + "douyin-creator-marketplace-xingtu": "douyin-xingtu", + "xiaohongshu-rednote": "xiaohongshu", + "xiaohongshu-creator-marketplace-pugongying": "xiaohongshu-pugongying", + "qq-huxuan-creator-marketplace": "qq-huxuan", + "wechat-official-accounts": "wechat-mp", + "taobao-and-tmall": "taobao", + "xianyu-goofish": "xianyu", + "jdcom": "jd", + "twitter": "x", + "social-media": "web", # one cross-platform Chinese-web search endpoint + "llm": "doubao", # the family holds a single endpoint, Doubao's web answer +} + + +def ingest_justoneapi(refresh: bool) -> tuple[Path, dict]: + sitemap = fetch("https://docs.justoneapi.com/en/sitemap.xml", "justoneapi_sitemap.xml", + refresh=refresh).decode() + pages = sorted({ + m.group(1) for m in re.finditer(r"https://docs\.justoneapi\.com/en/api/([a-z0-9\-]+/[a-z0-9\-]+)", sitemap) + }) + + def load(rel: str) -> tuple[str, dict | None]: + platform, endpoint = rel.split("/") + try: + blob = fetch( + f"https://docs.justoneapi.com/openapi/{platform}/{endpoint}-en.json", + f"justoneapi_{platform}__{endpoint}.json", refresh=refresh, + ) + return rel, json.loads(blob) + except Exception: + return rel, None + + todo = [p for p in pages if not p.endswith("/")] + print(f" fetching {len(todo)} per-endpoint spec(s)…", file=sys.stderr) + with ThreadPoolExecutor(max_workers=8) as pool: + specs = dict(pool.map(load, todo)) + + skip = core_routes("justoneapi") + # Prices are dashboard-only (no public price API); scripts/data/justoneapi_prices.json is a + # hand-exported snapshot of the dashboard's Pricing page — re-export + rerun to refresh. + prices_file = Path(__file__).parent / "data" / "justoneapi_prices.json" + joa_prices: dict[str, float] = {} + if prices_file.is_file(): + joa_prices = json.loads(prices_file.read_text()).get("prices", {}) + endpoints, missing = [], [] + for rel, spec in sorted(specs.items()): + if not spec or not spec.get("paths"): + missing.append(rel) + continue + family = rel.split("/")[0] + tag_platform = next( + (t.get("x-platform-id") for t in (spec.get("tags") or []) if t.get("x-platform-id")), None) + raw = (tag_platform or family).replace("_", "-") + platform = JOA_PLATFORM.get(family) or JOA_PLATFORM.get(raw) or raw + for path, item in spec["paths"].items(): + for method, op in item.items(): + if method.lower() not in ("get", "post", "put", "patch", "delete"): + continue + method = method.upper() + if (method, path) in skip: + continue + desc = clean(op.get("description") or op.get("summary") + or spec["info"].get("description") or "") + summary = re.split(r"(?<=[.!?])\s+", desc)[0] if desc else spec["info"].get("title", rel) + entry = { + "id": slug_id("justoneapi", path.replace("/api/", "")), + "tier": "extended", + "platform": platform, + "method": method, + "path": path, + "summary": summary[:400], + "docs_url": f"https://docs.justoneapi.com/en/api/{rel}", + } + # per-op `summary` is Just One API's short label ("Best Sellers"); the per-file + # info.title ("Amazon Best Sellers API (V1)") backs it up + title = short_name(clean(op.get("summary") or "") or spec["info"].get("title", ""), + summary) + if title: + entry["name"] = title + price = joa_prices.get(path) + if price is not None: + entry["cost"] = {"type": "per_success", "value": price, "currency": "CNY", + "note": "billed only on success (code 0); errors free"} + endpoints.append(entry) + + source = { + "method": "per-endpoint openapi files linked from the docs sitemap", + "ingested": str(date.today()), + "spec_urls": [ + "https://docs.justoneapi.com/en/sitemap.xml", + "https://docs.justoneapi.com/openapi//-en.json", + ], + } + notes = [ + "", + "Just One API publishes no combined spec; every documented endpoint has its own OpenAPI file", + "at /openapi//-en.json, enumerated here from the English sitemap.", + "`platform` comes from the spec's own x-platform-id tag, falling back to the docs URL segment.", + "Endpoints whose docs page carries a `-deprecated` slug are kept — the provider still serves", + "them — and are visible as such in the id. Prices come from scripts/data/justoneapi_prices.json,", + "a hand-exported snapshot of the logged-in dashboard's Pricing page (CNY per successful request);", + "an endpoint without a cost was absent from that export (e.g. not activated for the account).", + ] + if missing: + notes.append(f"{len(missing)} doc page(s) had no fetchable spec and were skipped.") + return write_extended("justoneapi", source, endpoints, notes), {"missing": len(missing)} + + +INGESTERS = {"tikhub": ingest_tikhub, "dataforseo": ingest_dataforseo, "justoneapi": ingest_justoneapi} + + +# ============================================================================================= +# FIRST-PARTY OAUTH PROVIDERS +# +# The scraper providers above sell breadth, and their extended tier reads as a menu. These nine +# are the opposite: one connected account, and the question the founder asked is "what can this +# credential actually DO?". Three things are therefore different here and are worth stating once: +# +# - **Scope gaps are data, not a filter.** A method whose Google/X scopes are NOT covered by what +# treg's OAuth apps request is still listed, carrying `scope_gap:` naming the missing scope. +# Dropping them would hide exactly the list someone needs in order to decide which scopes to +# add to the registered apps. Nothing generated here is callable-by-default in the way a +# scraper route is, so "listed" never implies "works today". +# - **Some methods live on a SIBLING HOST.** Google splits one product across several +# `*.googleapis.com` services (GA4 reporting vs GA4 admin; four separate My Business services) +# while `OAuthProvider.base_url` names exactly one. Those entries carry `host:` and their +# `path` is relative to THAT host, not to base_url. This is the same class of bug as the +# DataForSEO `/v3` note in docs/context/architecture/catalog.md — recorded explicitly instead +# of silently producing paths that 404 against the provisioned tool. +# - **No test_request, no verification.** These call a real business's own account; a generated +# test request would need a property id, a customer id or a Page id that no spec can supply. +# A later wave replays them through `--via-treg` with a live connection. +# +# Core-wins dedup compares paths with their placeholders NORMALISED ({property} == {property_id}), +# because a hand-curated core file and a machine spec never agree on placeholder spelling — the +# naive (method, path) comparison is what let every DataForSEO core route reappear in extended. + + +def _route_key(method: str, path: str) -> tuple[str, str]: + return str(method).upper(), re.sub(r"\{[^}]*\}", "{}", str(path)) + + +def core_route_keys(provider: str) -> set[tuple[str, str]]: + return {_route_key(m, p) for m, p in core_routes(provider)} + + +def _oauth_source(method: str, urls: list[str]) -> dict: + return {"method": method, "ingested": str(date.today()), "spec_urls": urls} + + +# --- Google discovery documents ---------------------------------------------------------------- +# Every Google REST API publishes a machine-readable Discovery document describing every method: +# httpMethod, flatPath (the URL relative to the service host), a description, the full parameter +# list with types/required flags/enums, and the OAuth scopes the method accepts. + +_DISCOVERY = "https://{svc}.googleapis.com/$discovery/rest?version={ver}" + + +def google_discovery(service: str, version: str, refresh: bool) -> dict: + return json.loads(fetch( + _DISCOVERY.format(svc=service, ver=version), + f"google_discovery_{service}_{version}.json", refresh=refresh, headers=_JSON_HDR, + )) + + +def google_methods(doc: dict) -> list[dict]: + """Every method in a discovery document, flattened out of its nested resource tree.""" + out: list[dict] = [] + + def walk(res: dict) -> None: + for _, r in sorted((res.get("resources") or {}).items()): + for _, m in sorted((r.get("methods") or {}).items()): + out.append(m) + walk(r) + + walk(doc) + return out + + +def _google_input(m: dict) -> dict: + """`input` from the discovery method's parameter block — names, types, required, enums.""" + inp: dict = {} + path_p: dict = {} + query_p: dict = {} + for name, p in sorted((m.get("parameters") or {}).items()): + entry: dict = {"type": p.get("type", "string"), "required": bool(p.get("required"))} + note = clean(p.get("description") or "") + if note: + entry["note"] = note[:300] + if p.get("enum"): + entry["enum"] = p["enum"] + (path_p if p.get("location") == "path" else query_p)[name] = entry + if path_p: + inp["pathParams"] = path_p + if query_p: + inp["queryParams"] = query_p + if m.get("request"): + ref = str((m["request"] or {}).get("$ref") or "object") + inp["bodyType"] = "json" + inp["body"] = { + "type": "object", "required": True, + "note": f"JSON body — the {ref} resource; see the method's API reference for its fields", + } + return inp + + +def _scope_gap(needed: list[str], granted: set[str]) -> str: + """Empty when this credential can call the method; otherwise what is missing, in one line. + + Google's discovery lists ALTERNATIVE scopes (holding any one of them suffices), so the test is + an intersection, not a subset. + """ + if not needed or set(needed) & granted: + return "" + short = [s.rsplit("/", 1)[-1] for s in needed] + return ("treg's OAuth app does not request this; needs one of: " + ", ".join(sorted(short))) + + +def google_entry( + provider: str, m: dict, *, platform: str, granted: set[str], + host: str = "", scope: str = "own_account", cost: dict | None = None, docs_url: str = "", +) -> dict: + """One extended entry from one discovery method. + + `path` prefers the media-upload URL when the method has one — an upload goes to + /upload//… and posting the metadata URL instead is a silent 400. + """ + upload = (((m.get("mediaUpload") or {}).get("protocols") or {}).get("simple") or {}).get("path") + path = upload or ("/" + str(m.get("flatPath") or m.get("path") or "").lstrip("/")) + mid = str(m["id"]) + desc = clean(m.get("description") or "") + tail = ".".join(mid.split(".")[-2:]) + entry: dict = { + "id": f"{provider}.x." + re.sub(r"[^a-z0-9]+", "-", mid.lower()).strip("-"), + "tier": "extended", + "platform": platform, + "method": str(m["httpMethod"]).upper(), + "path": path, + "summary": (f"{tail} — {desc}" if desc else tail)[:400], + "scope": scope, + } + if host: + entry["host"] = host + if cost: + entry["cost"] = cost + if docs_url: + entry["docs_url"] = docs_url + gap = _scope_gap(m.get("scopes") or [], granted) + if gap: + entry["scope_gap"] = gap + inp = _google_input(m) + if inp: + entry["input"] = inp + return entry + + +FREE_QUOTA = {"type": "free", "value": 0.0, "currency": "USD", + "note": "no per-call charge; billed against the API's daily quota"} + + +# --- google-search-console --------------------------------------------------------------------- + +def ingest_google_search_console(refresh: bool) -> tuple[Path, dict]: + provider = "google-search-console" + doc = google_discovery("searchconsole", "v1", refresh) + granted = { + "https://www.googleapis.com/auth/webmasters", + "https://www.googleapis.com/auth/webmasters.readonly", + } + skip = core_route_keys(provider) + endpoints, gaps = [], 0 + for m in google_methods(doc): + # the Mobile-Friendly Test declares no scope at all: it tests any public URL and merely + # happens to be hosted here, so it is not an own-account read + public = not (m.get("scopes") or []) + e = google_entry(provider, m, platform="search-console", granted=granted, cost=FREE_QUOTA, + scope="any_account" if public else "own_account", + docs_url="https://developers.google.com/webmaster-tools/v1/api_reference_index") + if _route_key(e["method"], e["path"]) in skip: + continue + gaps += bool(e.get("scope_gap")) + endpoints.append(e) + notes = [ + "", + "Source: the searchconsole v1 Discovery document, which covers BOTH surfaces of the product —", + "the legacy `webmasters/v3/…` resources (sites, sitemaps, searchanalytics) and the newer", + "`v1/…` ones (URL Inspection, the Mobile-Friendly Test). Both are served by", + "searchconsole.googleapis.com, so every path here is relative to the provider's base_url.", + "", + "treg requests webmasters.readonly by default and webmasters for the `write` capability, so", + "every method is reachable — sites.add/delete and sitemaps.submit/delete need `write`.", + "urlTestingTools.mobileFriendlyTest.run declares NO scope: it is a public tool that happens to", + "live on this host, so it is marked scope any_account.", + ] + return write_extended(provider, _oauth_source( + "google discovery document", + ["https://searchconsole.googleapis.com/$discovery/rest?version=v1"], + ), endpoints, notes), {"scope_gaps": gaps} + + +# --- google-analytics -------------------------------------------------------------------------- + +def ingest_google_analytics(refresh: bool) -> tuple[Path, dict]: + provider = "google-analytics" + granted = {"https://www.googleapis.com/auth/analytics.readonly"} + skip = core_route_keys(provider) + endpoints, gaps, admin = [], 0, 0 + for m in google_methods(google_discovery("analyticsdata", "v1beta", refresh)): + e = google_entry(provider, m, platform="google-analytics", granted=granted, cost=FREE_QUOTA, + docs_url="https://developers.google.com/analytics/devguides/reporting/data/v1") + if _route_key(e["method"], e["path"]) in skip: + continue + gaps += bool(e.get("scope_gap")) + endpoints.append(e) + for m in google_methods(google_discovery("analyticsadmin", "v1beta", refresh)): + e = google_entry(provider, m, platform="google-analytics", granted=granted, cost=FREE_QUOTA, + host="analyticsadmin.googleapis.com", + docs_url="https://developers.google.com/analytics/devguides/config/admin/v1") + if _route_key(e["method"], e["path"]) in skip: + continue + gaps += bool(e.get("scope_gap")) + admin += 1 + endpoints.append(e) + notes = [ + "", + "TWO HOSTS. Google splits GA4 in half: reporting (runReport, runRealtimeReport, funnels,", + "pivots) is analyticsdata.googleapis.com — the provider's base_url — while everything that", + "DESCRIBES the account (properties, data streams, custom dimensions, audiences, links to", + "Ads/BigQuery/Search Console) is analyticsadmin.googleapis.com.", + "", + f"The {admin} Admin API entries below carry `host: analyticsadmin.googleapis.com` and their", + "`path` is relative to THAT host. They are listed rather than dropped because the same OAuth", + "token calls both and they are most of what a connected GA4 account can do — but the tool", + "treg auto-provisions is bound to analyticsdata, so calling one needs a second tool bound to", + "the admin host (or `treg connections resources`, which already does the property listing).", + "Every entry WITHOUT a `host:` is callable through the provisioned google-analytics tool.", + "", + "treg requests analytics.readonly only. Every Admin write method (create/update/delete/", + "archive) therefore carries `scope_gap:` naming analytics.edit — that list is the answer to", + "'which scope would make GA4 writable'.", + ] + return write_extended(provider, _oauth_source( + "google discovery documents (analyticsdata + analyticsadmin)", + ["https://analyticsdata.googleapis.com/$discovery/rest?version=v1beta", + "https://analyticsadmin.googleapis.com/$discovery/rest?version=v1beta"], + ), endpoints, notes), {"scope_gaps": gaps, "admin": admin} + + +# --- google-business-profile ------------------------------------------------------------------- +# Google retired the single "My Business API v4" into SIX narrow services, each on its own host. +# base_url is the account-management one, so only that family is callable through the provisioned +# tool; the rest carry `host:`. Reviews never made the migration and still live on the legacy v4 +# host, which publishes NO discovery document — those few routes are hand-listed and flagged. + +GBP_SERVICES = [ + ("mybusinessaccountmanagement", "v1", ""), # == base_url: no host override needed + ("mybusinessbusinessinformation", "v1", "mybusinessbusinessinformation.googleapis.com"), + ("businessprofileperformance", "v1", "businessprofileperformance.googleapis.com"), + ("mybusinessqanda", "v1", "mybusinessqanda.googleapis.com"), + ("mybusinessplaceactions", "v1", "mybusinessplaceactions.googleapis.com"), + ("mybusinessnotifications", "v1", "mybusinessnotifications.googleapis.com"), + ("mybusinessverifications", "v1", "mybusinessverifications.googleapis.com"), +] + +# The legacy v4 surface, which has no discovery document (https://mybusiness.googleapis.com/ +# $discovery/rest?version=v4 answers 404). Reviews are the reason anyone connects this provider, +# so the handful that matter are transcribed from the published v4 reference. Kept deliberately +# short — hand-transcribed paths are how typos ship (catalog.md, step 1). +GBP_V4 = [ + ("reviews-get", "GET", "/v4/accounts/{account_id}/locations/{location_id}/reviews/{review_id}", + "reviews.get — Read a single review left on one of your listings"), + ("reviews-deletereply", "DELETE", + "/v4/accounts/{account_id}/locations/{location_id}/reviews/{review_id}/reply", + "reviews.deleteReply — Delete the business's public reply to a review"), + ("reviews-batchget", "GET", "/v4/accounts/{account_id}/locations:batchGetReviews", + "reviews.batchGetReviews — Read reviews across several of your listings in one call"), + ("media-list", "GET", "/v4/accounts/{account_id}/locations/{location_id}/media", + "media.list — List the photos and videos on one of your listings"), + ("media-create", "POST", "/v4/accounts/{account_id}/locations/{location_id}/media", + "media.create — Add a photo or video to one of your listings"), + ("localposts-list", "GET", "/v4/accounts/{account_id}/locations/{location_id}/localPosts", + "localPosts.list — List the local posts (offers, updates, events) on one of your listings"), + ("localposts-create", "POST", "/v4/accounts/{account_id}/locations/{location_id}/localPosts", + "localPosts.create — Publish a local post (offer, update or event) to one of your listings"), +] + + +def ingest_google_business_profile(refresh: bool) -> tuple[Path, dict]: + provider = "google-business-profile" + granted = {"https://www.googleapis.com/auth/business.manage"} + skip = core_route_keys(provider) + endpoints, off_host = [], 0 + for service, version, host in GBP_SERVICES: + for m in google_methods(google_discovery(service, version, refresh)): + e = google_entry(provider, m, platform="google-business", granted=granted, + host=host, cost=FREE_QUOTA, + docs_url=f"https://developers.google.com/my-business/reference/{service}/rest") + if _route_key(e["method"], e["path"]) in skip: + continue + # discovery ids collide across the six services (accounts.list exists in three of them) + e["id"] = f"{provider}.x.{service.removeprefix('mybusiness')}-" + e["id"].split(".x.", 1)[1] + off_host += bool(host) + endpoints.append(e) + for slug, method, path, summary in GBP_V4: + if _route_key(method, path) in skip: + continue + endpoints.append({ + "id": f"{provider}.x.v4-{slug}", "tier": "extended", "platform": "google-business", + "method": method, "path": path, "summary": summary, "scope": "own_account", + "host": "mybusiness.googleapis.com", "cost": FREE_QUOTA, + "docs_url": "https://developers.google.com/my-business/reference/rest/v4/accounts.locations.reviews", + }) + notes = [ + "", + "SIX HOSTS, plus a legacy one. Google broke the old My Business API into narrow services, one", + "host each. base_url is mybusinessaccountmanagement.googleapis.com, so only that family is", + "callable through the provisioned tool; every other entry carries `host:` and its `path` is", + "relative to THAT host. Services covered: accountmanagement (accounts, admins, invitations),", + "businessinformation (locations, attributes, categories, chains), businessprofileperformance", + "(daily metrics and search keywords), qanda (questions and answers), placeactions, ", + "notifications, verifications.", + "", + "Reviews never migrated: they are still on mybusiness.googleapis.com/v4, which publishes NO", + "discovery document. The 7 v4 entries here are hand-transcribed from the published reference", + "and are the only ones in this file not machine-derived — treat their paths with more", + "suspicion than the rest. (The core file's two review endpoints have the same origin.)", + "", + "No scope_gap anywhere: the My Business discovery documents declare no scopes at all, so", + "coverage cannot be computed from the spec. business.manage — what treg requests — is in", + "practice the single scope the whole family uses. The real gate on this provider is not the", + "scope but Google's separate Business Profile API quota grant, which starts every project at", + "zero requests per day.", + ] + return write_extended(provider, _oauth_source( + "google discovery documents (six My Business services) + hand-listed legacy v4 reviews", + [_DISCOVERY.format(svc=s, ver=v) for s, v, _ in GBP_SERVICES] + + ["https://developers.google.com/my-business/reference/rest/v4/accounts.locations.reviews"], + ), endpoints, notes), {"off_host": off_host} + + +# --- youtube ----------------------------------------------------------------------------------- +# Quota, not money, is what runs out on the YouTube Data API: a project gets 10,000 units a day and +# a single search.list costs 100 of them. The per-method costs below are the published quota table +# (developers.google.com/youtube/v3/determine_quota_cost) — the discovery document does not carry +# them, and an agent that does not know search.list is 100x a videos.list will exhaust a day's +# quota in 100 calls. +YT_QUOTA = { + "youtube.search.list": 100, + "youtube.captions.list": 50, "youtube.captions.insert": 400, "youtube.captions.update": 450, + "youtube.captions.download": 200, "youtube.captions.delete": 50, + "youtube.videos.insert": 1600, "youtube.videos.getRating": 1, +} +# Reads that return the same data for anybody — the connected account is just how we authenticate. +YT_PUBLIC = { + "youtube.videos.list", "youtube.search.list", "youtube.channels.list", + "youtube.commentThreads.list", "youtube.comments.list", "youtube.playlists.list", + "youtube.playlistItems.list", "youtube.channelSections.list", "youtube.activities.list", + "youtube.videoCategories.list", "youtube.guideCategories.list", "youtube.i18nLanguages.list", + "youtube.i18nRegions.list", "youtube.playlistImages.list", "youtube.members.list", +} + + +def _yt_quota(mid: str) -> int: + if mid in YT_QUOTA: + return YT_QUOTA[mid] + return 1 if mid.endswith(".list") else 50 + + +def ingest_youtube(refresh: bool) -> tuple[Path, dict]: + provider = "youtube" + doc = google_discovery("youtube", "v3", refresh) + granted = { + "https://www.googleapis.com/auth/youtube.readonly", + "https://www.googleapis.com/auth/youtube.upload", + "https://www.googleapis.com/auth/youtube", + "https://www.googleapis.com/auth/youtube.force-ssl", + } + skip = core_route_keys(provider) + endpoints, gaps, units = [], 0, 0 + for m in google_methods(doc): + mid = str(m["id"]) + q = _yt_quota(mid) + cost = {"type": "free", "value": 0.0, "currency": "USD", + "note": f"{q} quota unit{'s' if q != 1 else ''} per call " + f"(default project allowance: 10,000 units/day)"} + e = google_entry(provider, m, platform="youtube", granted=granted, cost=cost, + scope="any_account" if mid in YT_PUBLIC else "own_account", + docs_url="https://developers.google.com/youtube/v3/docs") + if _route_key(e["method"], e["path"]) in skip: + continue + gaps += bool(e.get("scope_gap")) + units += q + endpoints.append(e) + notes = [ + "", + "Source: the youtube v3 Discovery document — the complete Data API. Paths are relative to", + "base_url (youtube.googleapis.com); the three upload methods correctly point at the", + "/upload/… URL rather than the metadata one, which is a 400 if you post a body to it.", + "", + "COST IS QUOTA, NOT MONEY. Nothing here is billed, but a project gets 10,000 units a day and", + "the methods are not equal: a list is 1 unit, a write is 50, search.list is 100, captions", + "cost 200-450 and a video UPLOAD is 1600 — six uploads is more than half a day's allowance.", + "Each entry's cost.note carries its published unit price (from YouTube's quota-cost table;", + "the discovery document does not include it).", + "", + "scope: `any_account` marks the reads that return the same public data for anyone (videos,", + "search, channels, comments, playlists) — the connection is just how the call authenticates.", + "Everything else acts on the connected channel.", + "", + "Only 2 scope gaps, both channel memberships (members.list, membershipsLevels.list), which", + "need youtube.channel-memberships.creator. The `manage` capability's four scopes cover the", + "entire rest of the Data API — including the content-owner methods, which accept plain", + "`youtube` as an alternative to youtubepartner and so are NOT gaps, though they are only", + "meaningful to a CMS partner account.", + ] + return write_extended(provider, _oauth_source( + "google discovery document", + ["https://youtube.googleapis.com/$discovery/rest?version=v3", + "https://developers.google.com/youtube/v3/determine_quota_cost"], + ), endpoints, notes), {"scope_gaps": gaps, "quota_units": units} + + +# --- google-ads -------------------------------------------------------------------------------- +# The Google Ads API has no discovery document or OpenAPI spec worth ingesting: apart from a +# handful of mutate services, ONE endpoint (googleAds:searchStream) answers everything, and what +# varies is the GAQL `FROM` clause. So the unit of coverage here is the RESOURCE, not the route — +# forty entries pointing at the same path, each documenting one thing you can query. That is what +# an agent actually needs: it already knows how to POST a query, it does not know that +# `search_term_view` is where the actual user searches live. +GADS_VERSION = "v21" # must track OAuthProvider.examples / core google-ads.yaml +GADS_RESOURCES: list[tuple[str, str]] = [ + ("campaign", "Campaigns — name, status, budget, bidding strategy, start/end dates and all campaign-level metrics"), + ("ad_group", "Ad groups — name, status, type, CPC bid and ad-group-level metrics"), + ("ad_group_ad", "The individual ads — headlines, descriptions, final URLs, approval status and per-ad metrics"), + ("ad_group_criterion", "Everything targeted or excluded in an ad group: keywords, audiences, placements, with bids"), + ("keyword_view", "Per-keyword performance — the metrics view over ad_group_criterion keywords"), + ("search_term_view", "The actual queries people typed that triggered your ads, and what each cost"), + ("campaign_budget", "Budgets — daily amount, delivery method, and which campaigns share them"), + ("campaign_criterion", "Campaign-level targeting and exclusions: locations, languages, devices, negative keywords"), + ("customer", "The account itself — name, currency, timezone, auto-tagging, and account-level metrics"), + ("customer_client", "Every account under a manager (MCC), including the hierarchy level and whether it is a manager"), + ("conversion_action", "Conversion actions — what counts as a conversion, its category, value and counting rules"), + ("change_event", "Who changed what in the account in the last 30 days, old value and new value"), + ("change_status", "Which resources changed since a timestamp — the cheap way to poll for edits"), + ("bidding_strategy", "Portfolio bidding strategies and their targets (target CPA, target ROAS, maximise conversions)"), + ("billing_setup", "Billing setups and payment accounts attached to the customer"), + ("account_budget", "Account-level budgets and their approved spending limits (invoiced accounts)"), + ("asset", "Assets — images, videos, text, sitelinks, callouts — and where each is used"), + ("asset_group", "Asset groups (Performance Max) with their status and per-group metrics"), + ("campaign_asset", "Which assets are attached to which campaign, and in what field type"), + ("customer_asset", "Assets attached at the account level"), + ("ad_group_ad_asset_view", "Per-asset performance inside a specific ad — which headline actually pulled"), + ("audience", "Audiences defined in the account and their segments"), + ("user_list", "Remarketing and customer-match lists, their size and eligibility per network"), + ("user_interest", "The affinity and in-market interest taxonomy available for targeting"), + ("campaign_audience_view", "Audience performance at campaign level"), + ("ad_group_audience_view", "Audience performance at ad-group level"), + ("age_range_view", "Performance split by the viewer's age bracket"), + ("gender_view", "Performance split by the viewer's gender"), + ("parental_status_view", "Performance split by parental status"), + ("income_range_view", "Performance split by household income bracket"), + ("geographic_view", "Performance by the location that triggered the ad (physical or interest)"), + ("user_location_view", "Performance by where the user actually was, ignoring location-of-interest"), + ("location_view", "Performance per targeted location criterion"), + ("ad_schedule_view", "Performance by day of week and hour of day"), + ("detail_placement_view", "Exact URLs and apps your Display/Video ads ran on"), + ("managed_placement_view", "Performance of placements you chose yourself"), + ("topic_view", "Performance of Display topic targeting"), + ("landing_page_view", "Performance grouped by the unexpanded final URL"), + ("shopping_performance_view", "Shopping metrics broken down by product attributes (brand, category, item id)"), + ("recommendation", "Google's own optimisation suggestions for the account, with their estimated impact"), + ("paid_organic_search_term_view", "Paid and organic side by side per query — needs Search Console linked"), + ("click_view", "Individual clicks with GCLID and location, for the last 90 days"), +] + + +def ingest_google_ads(refresh: bool) -> tuple[Path, dict]: + provider = "google-ads" + skip = core_route_keys(provider) + path = f"/{GADS_VERSION}/customers/{{customer_id}}/googleAds:searchStream" + endpoints = [] + for resource, summary in GADS_RESOURCES: + if _route_key("POST", path) in skip: + break + endpoints.append({ + "id": f"{provider}.x.gaql-{resource.replace('_', '-')}", + "tier": "extended", + "platform": "google-ads", + "method": "POST", + "path": path, + "summary": f"GAQL FROM {resource} — {summary}", + "scope": "own_account", + "cost": {"type": "free", "value": 0.0, "currency": "USD", + "note": "no per-call charge; counts against the developer token's daily " + "operation limit (Basic access: 15,000 operations/day)"}, + "input": { + "pathParams": {"customer_id": { + "type": "string", "required": True, + "note": "the 10-digit Google Ads customer id, no dashes; from listAccessibleCustomers"}}, + "bodyType": "json", + "body": {"type": "object", "required": True, + "note": f'{{"query": "SELECT FROM {resource} ' + f'WHERE segments.date DURING LAST_30_DAYS"}}'}, + "note": f"Resource: {resource}. Its selectable fields, segments and metrics are listed at " + f"https://developers.google.com/google-ads/api/fields/{GADS_VERSION}/{resource}. " + f"searchStream returns the whole result set in one streamed response and takes no " + f"page size; googleAds:search is the paginated twin (curated in the core file).", + }, + "docs_url": f"https://developers.google.com/google-ads/api/fields/{GADS_VERSION}/{resource}", + }) + notes = [ + "", + "NOT A ROUTE LIST — A RESOURCE LIST. The Google Ads API publishes no discovery document or", + "OpenAPI spec, and it would not help if it did: one endpoint, googleAds:searchStream, answers", + "every read, and what actually varies is the GAQL `FROM` clause. So each entry below is one", + f"QUERYABLE RESOURCE pointed at that single path, capped at the {len(GADS_RESOURCES)} an agent would realistically", + "reach for out of ~180 the API defines. `input.note` names the resource and links its field", + "reference; `docs_url` goes straight to the selectable fields, segments and metrics.", + "", + f"API version {GADS_VERSION}, matching the core file and OAuthProvider.examples. A wrong version 404s as", + "HTML rather than JSON, so this must be bumped in all three places together.", + "", + "Every call needs BOTH the OAuth bearer and a developer-token header from an approved manager", + "(MCC) account — see OAuthProvider.extra_credential_note. The adwords scope treg requests", + "covers the whole API, so there are no scope gaps here; the gate is the developer token's", + "access level (Basic caps the account at 15,000 operations/day).", + ] + return write_extended(provider, _oauth_source( + "official GAQL resource reference (one entry per queryable resource)", + [f"https://developers.google.com/google-ads/api/fields/{GADS_VERSION}/overview", + f"https://developers.google.com/google-ads/api/reference/rpc/{GADS_VERSION}/overview"], + ), endpoints, notes), {} + + +# --- x ------------------------------------------------------------------------------------------ + +X_OPENAPI = "https://api.twitter.com/2/openapi.json" +# treg's X app requests these across its read + write capabilities. Anything else is a scope gap. +X_GRANTED = {"tweet.read", "tweet.write", "users.read", "offline.access"} +# Path fragments that make a route act on the CONNECTED account rather than read public data. +X_OWN = re.compile( + r"/(me|dm_conversations|dm_events|bookmarks|blocking|muting|compliance|usage" + r"|account_activity|webhooks|reverse_chronological)\b") + + +def _x_scopes(op: dict) -> list[str]: + for sec in op.get("security") or []: + for name, scopes in sec.items(): + if name.startswith("OAuth2User") and scopes: + return list(scopes) + return [] + + +def ingest_x(refresh: bool) -> tuple[Path, dict]: + provider = "x" + spec = json.loads(fetch(X_OPENAPI, "x_openapi.json", refresh=refresh, headers=_JSON_HDR)) + skip = core_route_keys(provider) + endpoints, gaps = [], 0 + for path, item in sorted((spec.get("paths") or {}).items()): + for method, op in item.items(): + if method.lower() not in ("get", "post", "put", "patch", "delete"): + continue + method = method.upper() + if _route_key(method, path) in skip: + continue + opid = str(op.get("operationId") or "") + slug = re.sub(r"(? tuple[Path, dict]: + platform, edges, granted, docs = META_PROVIDERS[service] + # EXACT (method, path) dedup here, not the normalised one the Google ingesters use. On Graph the + # node id IS the first path segment, so `/{post_id}/insights` and `/{page_id}/insights` differ + # only by the placeholder NAME while being genuinely different endpoints — normalising would + # silently drop post insights because the core file curates page insights. The paths below are + # hand-written against the core files, so exact matching is enough. + skip = core_routes(service) + endpoints, gaps = [], 0 + for slug, method, path, summary, gap in edges: + if _route_key(method, path) in skip: + continue + entry: dict = { + "id": f"{service}.x.{slug}", "tier": "extended", "platform": platform, + "method": method, "path": path, "summary": summary, "scope": "own_account", + "cost": {"type": "free", "value": 0.0, "currency": "USD", + "note": "no per-call charge; counted against the app's Graph API rate limit"}, + "docs_url": docs, + } + if gap: + gaps += 1 + entry["scope_gap"] = f"treg's OAuth app does not request this; {gap}" + endpoints.append(entry) + notes = [ + "", + "META PUBLISHES NO SPEC. There is no OpenAPI document, no discovery document and no", + "machine-readable index for the Graph API — the only source is the HTML reference, one page", + "per node type. Hand-transcribing hundreds of routes out of HTML is how typos ship, so this", + f"file is deliberately SHORT: {len(endpoints)} entries, each an edge of a node type this provider's scopes", + "can actually reach, read off the official reference page linked in every docs_url. Accuracy", + "over volume — assume the Graph surface is much larger than what is listed here.", + "", + "Paths are relative to base_url, which ALREADY ENDS IN /v25.0 — no entry repeats the version.", + "Bumping the Graph version is a base_url change in oauth_providers.py, not an edit here.", + "", + f"treg's app requests {granted}.", + f"{gaps} entr{'ies' if gaps != 1 else 'y'} need a scope beyond that and carry `scope_gap:`. Meta gates most of those behind", + "App Review with a screencast per permission, so a scope gap here is a product decision, not", + "a config change.", + "", + "Graph ids are opaque and account-specific ({page_id}, {ad_account_id} — the latter includes", + "its `act_` prefix), so nothing here carries a test_request; `treg connections resources`", + "supplies the id for the connected account.", + ] + return write_extended(service, _oauth_source( + "hand-curated from the official Graph API HTML reference (Meta publishes no machine-readable spec)", + [docs, META_GRAPH_DOCS], + ), endpoints, notes), {"scope_gaps": gaps} + + +INGESTERS.update({ + "google-search-console": ingest_google_search_console, + "google-analytics": ingest_google_analytics, + "google-business-profile": ingest_google_business_profile, + "youtube": ingest_youtube, + "google-ads": ingest_google_ads, + "x": ingest_x, + "facebook": lambda refresh: ingest_meta("facebook", refresh), + "instagram": lambda refresh: ingest_meta("instagram", refresh), + "meta-ads": lambda refresh: ingest_meta("meta-ads", refresh), +}) + + +def main(argv: list[str]) -> int: + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("providers", nargs="+", choices=[*INGESTERS, "all"]) + ap.add_argument("--refresh", action="store_true", help="re-download instead of using the cache") + args = ap.parse_args(argv) + + names = list(INGESTERS) if "all" in args.providers else args.providers + platforms = set(yaml.safe_load((CATALOG / "capabilities.yaml").read_text())["platforms"]) + rc = 0 + for name in names: + print(f"{name}:", file=sys.stderr) + path, _ = INGESTERS[name](args.refresh) + data = yaml.safe_load(path.read_text()) + used = {} + for ep in data["endpoints"]: + used[ep["platform"]] = used.get(ep["platform"], 0) + 1 + unknown = sorted(set(used) - platforms) + print(f" {len(data['endpoints'])} endpoint(s) → {path.relative_to(ROOT)}") + print(" platforms: " + ", ".join(f"{k} ({v})" for k, v in sorted(used.items(), key=lambda kv: -kv[1]))) + if unknown: + rc = 1 + print(f" MISSING from capabilities.yaml platforms: {', '.join(unknown)}", file=sys.stderr) + return rc + + +if __name__ == "__main__": + raise SystemExit(main(sys.argv[1:])) diff --git a/scripts/catalog_validate.py b/scripts/catalog_validate.py new file mode 100644 index 00000000..a78e9f30 --- /dev/null +++ b/scripts/catalog_validate.py @@ -0,0 +1,227 @@ +#!/usr/bin/env python3 +"""Validate the endpoint catalog (src/treg/catalog/*.yaml) — schema shape + referential integrity. + +Run from the repo root: `uv run python scripts/catalog_validate.py [service ...]` +Exit 0 = valid. Every violation prints one line: `: `. + +Checks (the success criteria from docs/context/architecture/catalog.md): + - provider file's `provider` matches its filename and exists in treg.oauth_providers.REGISTRY + - endpoint ids unique across the WHOLE catalog; id convention `.` + - `capability` exists in capabilities.yaml OR the file's own proposed_capabilities + - `platform` equals the capability's first segment and exists in capabilities.yaml platforms + - required fields present; enums valid (scope, method, cost.type) + - a `verified` endpoint must have an existing example_response file + - an extended endpoint a verification run has touched claims exactly one non-empty state + (verified | unverified | untestable | skipped), and an `untestable` one carries no + test_request for a re-verify run to call it with anyway + - no obvious credential leak (Authorization/token values) in any catalog file + +Two tiers, two rule sets. `.yaml` holds hand-curated `tier: core` entries and every rule +above applies. `.extended.yaml` holds the machine-generated coverage tier written by +scripts/catalog_ingest.py: those entries need only id/platform/method/path/summary, and carry no +capability (nothing has been mapped to the taxonomy yet). Id-prefix, id-uniqueness, platform +integrity and the leak scan apply to BOTH — an extended entry that does declare a capability is +held to the core referential rules. + +The evidence rule is deliberately tier-blind, and needed no loosening when the extended tier gained +generated `test_request`s and live verification (scripts/catalog_verify_extended.py): `verified` +means the same thing in both files, so in both it must be backed by an `example_response` file that +exists and a `test_request` to re-check it with. A tier decides how much is EXPECTED of an entry, +never what a claim on it is worth. +""" + +from __future__ import annotations + +import re +import sys +from pathlib import Path + +import yaml + +ROOT = Path(__file__).resolve().parent.parent +CATALOG = ROOT / "src" / "treg" / "catalog" + +SCOPES = {"any_account", "own_account"} +METHODS = {"GET", "POST", "PUT", "PATCH", "DELETE"} +COST_TYPES = {"per_call", "per_result", "per_success", "free", "quota_rows"} +TIERS = {"core", "extended"} +# the section heading an endpoint files under on its platform page — one lowercase word +DOMAIN = re.compile(r"[a-z][a-z0-9_]*") +REQUIRED = { + "core": ("id", "capability", "platform", "method", "path", "summary"), + "extended": ("id", "platform", "method", "path", "summary"), +} +# The four outcomes an extended endpoint can have once a verification run has touched it. They are +# mutually exclusive by construction, and an entry that has been through the pipeline must claim +# exactly one — "no state" reads as "never attempted", which is a different fact and a lie once a +# run has been over it. This caught a real regression: re-running an endpoint overwrote its result +# record, dropped the reason string, and stamped an EMPTY state that nothing else noticed. +STATES = ("verified", "unverified", "untestable", "skipped") +# a long token-looking literal anywhere in a catalog file is a leak until proven otherwise +LEAK = re.compile(r"(Bearer\s+[A-Za-z0-9+/_=-]{16,}|[A-Za-z0-9+/]{40,}={0,2})") +# ...but URL and API paths are also long runs of [A-Za-z0-9/], and the extended tier is thousands +# of them. Two things separate them from a credential: a path is built of short slash-separated +# segments, and those segments spell words ("dataforseo", "kolContentTags"). A base64 secret hits +# a `/` only about once per 64 characters, so at least one of its segments stays long and wordless. +WORDY = re.compile(r"[a-z]{8,}") + + +def looks_like_secret(match: str) -> bool: + if match.lower().startswith("bearer"): + return True + return any( + len(seg) >= 24 and not WORDY.search(seg) and any(c.isdigit() for c in seg) + for seg in match.split("/") + ) + + + +def fail(errors: list[str], where: str, msg: str) -> None: + errors.append(f"{where}: {msg}") + + +def main(argv: list[str]) -> int: + errors: list[str] = [] + tax = yaml.safe_load((CATALOG / "capabilities.yaml").read_text()) + platforms = set(tax.get("platforms") or {}) + capabilities = set(tax.get("capabilities") or {}) + + sys.path.insert(0, str(ROOT / "src")) + from treg.oauth_providers import REGISTRY # noqa: E402 + + only = set(argv) + files = sorted(p for p in CATALOG.glob("*.yaml") if p.name not in ("capabilities.yaml", "fx.yaml")) + # "tikhub" selects tikhub.yaml AND tikhub.extended.yaml — a service is both its tiers + service_of = {p: p.stem.removesuffix(".extended") for p in files} + if only: + files = [p for p in files if service_of[p] in only] + missing = only - {service_of[p] for p in files} + for m in missing: + fail(errors, m, "no such catalog file") + + seen_ids: dict[str, str] = {} + for path in files: + name = path.name + text = path.read_text() + data = yaml.safe_load(text) + if not isinstance(data, dict): + fail(errors, name, "not a mapping") + continue + extended_file = path.name.endswith(".extended.yaml") + service = data.get("provider") + if service != service_of[path]: + fail(errors, name, f"provider '{service}' != filename stem '{service_of[path]}'") + if service not in REGISTRY: + fail(errors, name, f"provider '{service}' not in oauth_providers.REGISTRY") + src = data.get("source") + # core files cite the docs they were curated from; extended files cite the specs they were + # generated from, so that a re-run is reproducible from the file alone + need = "spec_urls" if extended_file else "docs" + if not isinstance(src, dict) or not src.get(need): + fail(errors, name, f"source.{need} missing") + proposed = set(data.get("proposed_capabilities") or {}) + known = capabilities | proposed + + eps = data.get("endpoints") + if not isinstance(eps, list) or not eps: + fail(errors, name, "endpoints missing or empty") + continue + for ep in eps: + eid = ep.get("id", "") + where = f"{name}:{eid}" + tier = ep.get("tier", "extended" if extended_file else "core") + if tier not in TIERS: + fail(errors, where, f"bad tier '{tier}'") + tier = "extended" if extended_file else "core" + if extended_file != (tier == "extended"): + fail(errors, where, f"tier '{tier}' does not belong in {name}") + for f in REQUIRED[tier]: + if not ep.get(f): + fail(errors, where, f"missing required field '{f}'") + if eid in seen_ids: + fail(errors, where, f"duplicate id (also in {seen_ids[eid]})") + seen_ids[eid] = name + if service and not eid.startswith(f"{service}."): + fail(errors, where, f"id must start with '{service}.'") + # extended entries are unmapped by design; one that DOES claim a capability is held to + # the same referential rules as core, so a hand-promoted entry can't drift + cap = ep.get("capability", "") + if cap or tier == "core": + if cap not in known: + fail(errors, where, f"capability '{cap}' not in capabilities.yaml or proposed_capabilities") + plat = ep.get("platform", "") + if plat not in platforms: + fail(errors, where, f"platform '{plat}' not in capabilities.yaml platforms") + if cap and plat and cap.split(".")[0] != plat: + fail(errors, where, f"platform '{plat}' != capability's first segment '{cap.split('.')[0]}'") + # `domain` is optional — the loader derives one from the capability id or the path when + # it is absent. Declaring one overrides that, so it has to be the same SHAPE the derived + # ones are: one lowercase word, or the platform page grows a section of one. + # `name` is an optional short DISPLAY title; `summary` stays the provider's own + # description. Light check only: present ⇒ a non-empty string that fits a row heading. + nm = ep.get("name") + if nm is not None and (not isinstance(nm, str) or not nm.strip() or len(nm) > 60): + fail(errors, where, "name must be a non-empty string of at most 60 chars") + dom = ep.get("domain") + if dom is not None and not DOMAIN.fullmatch(str(dom)): + fail(errors, where, f"domain '{dom}' must be a single lowercase word (a-z0-9_)") + if ep.get("scope", "any_account") not in SCOPES: + fail(errors, where, f"bad scope '{ep.get('scope')}'") + if ep.get("method") not in METHODS: + fail(errors, where, f"bad method '{ep.get('method')}'") + cost = ep.get("cost") + if cost is not None or tier == "core": + # cost is optional in the extended tier — several providers publish prices per API + # family rather than per route — but a stated cost must still be a real cost model + if not isinstance(cost, dict) or cost.get("type") not in COST_TYPES: + fail(errors, where, f"cost.type missing or not one of {sorted(COST_TYPES)}") + if ep.get("verified"): + ex = ep.get("example_response") + if not ex: + fail(errors, where, "verified but no example_response") + elif not (CATALOG / ex).is_file(): + fail(errors, where, f"example_response '{ex}' does not exist") + if not ep.get("test_request"): + fail(errors, where, "verified but no test_request (nothing to re-verify with)") + if tier == "extended": + # "been through the pipeline" = a run either built it a request or recorded an + # outcome for it. A freshly ingested entry that has never been verified has + # neither, and is left alone. + claimed = [s for s in STATES if ep.get(s)] + touched = ep.get("test_request") or any(s in ep for s in STATES) + if len(claimed) > 1: + fail(errors, where, f"claims {len(claimed)} states at once: {claimed} — " + "verified/unverified/untestable/skipped are exclusive") + elif touched and not claimed: + empty = [s for s in STATES if s in ep] + fail(errors, where, f"no endpoint state: {sorted(empty)} present but empty" + if empty else "has a test_request but no verified/unverified/untestable/" + "skipped saying what happened when it was called") + if ep.get("untestable") and ep.get("test_request"): + # `untestable` means no call is possible — but catalog_verify.py --extended + # replays anything that HAS a test_request, so the pair is not just a + # contradiction on paper: it gets the endpoint called, and billed, by a + # re-verification run that was told it was uncallable. + fail(errors, where, "untestable but carries a test_request — a re-verify run " + "would call it anyway; drop one of the two") + + if extended_file: + # extended files are machine-generated from PUBLIC specs and public target ids (TikTok + # secUids, WeChat export ids, CDN URIs, pagination cursors) — long opaque strings that + # pattern-match as secrets endlessly. Credentials only ever travel via TREG_CATALOG_CRED + # env in the verify scripts, so here only the unambiguous leak shape is flagged. + for m in re.finditer(r"Bearer\s+[A-Za-z0-9+/_=-]{16,}", text): + fail(errors, name, f"credential literal in file: '{m.group(0)[:24]}…'") + continue + for m in LEAK.finditer(text): + if looks_like_secret(m.group(0)): + fail(errors, name, f"possible credential literal in file: '{m.group(0)[:24]}…'") + + for e in errors: + print(e) + print(f"{'FAIL' if errors else 'OK'} — {len(files)} provider file(s), {len(seen_ids)} endpoint(s), {len(errors)} error(s)") + return 1 if errors else 0 + + +if __name__ == "__main__": + raise SystemExit(main(sys.argv[1:])) diff --git a/scripts/catalog_verify.py b/scripts/catalog_verify.py new file mode 100644 index 00000000..dba8c951 --- /dev/null +++ b/scripts/catalog_verify.py @@ -0,0 +1,204 @@ +#!/usr/bin/env python3 +"""Live-verify catalog endpoints and capture example responses. + +Usage (credential comes ONLY from the environment — never from a file or argv): + TREG_CATALOG_CRED='' uv run python scripts/catalog_verify.py [--id ] [--write] + TREG_CATALOG_CRED='' uv run python scripts/catalog_verify.py --extended [--id …] + +For each endpoint in src/treg/catalog/.yaml it sends `test_request` to base_url+path with +the provider's auth shape (taken from treg.oauth_providers — token_header/format/location/encode), +checks `expect` (default: HTTP 2xx), and with --write saves a truncated example response to +src/treg/catalog/examples/.json. Prints one PASS/FAIL line per endpoint. + +`--extended` reads .extended.yaml instead, to REPLAY what a bulk verification run already +stamped there — it is a re-verifier, not the bulk runner (that is catalog_verify_extended.py). Three +differences the extended tier forces, all no-ops for core files: + * an entry without a `test_request` is skipped rather than called with no parameters; + * a leading duplicate of base_url's own path is stripped from `path`, because the extended tier + stores routes exactly as the provider's spec spells them and DataForSEO's include the /v3 that + base_url already ends with; + * `test_request.bodyType: form` sends the body as a form AND moves a query-param credential into + it — Just One API's POST routes read the token from the body and reject it in the query. + +It does NOT stamp `verified:` in the YAML — the curator does that, only for PASSes, so a stamp is +always a human-reviewed claim. Truncation: arrays -> first 2 items, strings -> 500 chars, and a +whole-document cap; the curator must still read every example for PII before committing. +""" + +from __future__ import annotations + +import argparse +import base64 +import json +import os +import sys +from datetime import date +from pathlib import Path +from urllib.parse import urlsplit + +import httpx +import yaml + +ROOT = Path(__file__).resolve().parent.parent +CATALOG = ROOT / "src" / "treg" / "catalog" +MAX_STR = 500 +MAX_LIST = 2 +MAX_BYTES = 10_000 + + +def truncate(node, depth=0): + if isinstance(node, dict): + return {k: truncate(v, depth + 1) for k, v in node.items()} + if isinstance(node, list): + out = [truncate(v, depth + 1) for v in node[:MAX_LIST]] + if len(node) > MAX_LIST: + out.append(f"… {len(node) - MAX_LIST} more item(s) truncated") + return out + if isinstance(node, str) and len(node) > MAX_STR: + return node[:MAX_STR] + f"… [{len(node)} chars total]" + return node + + +def _shrink_strings(node, limit: int): + if isinstance(node, dict): + return {k: _shrink_strings(v, limit) for k, v in node.items()} + if isinstance(node, list): + return [_shrink_strings(v, limit) for v in node] + if isinstance(node, str) and len(node) > limit: + return node[:limit] + "…" + return node + + +def dig(doc, dotted: str): + cur = doc + for part in dotted.split("."): + if isinstance(cur, list): + cur = cur[int(part)] + elif isinstance(cur, dict): + cur = cur.get(part) + else: + return None + return cur + + +def main() -> int: + ap = argparse.ArgumentParser() + ap.add_argument("service") + ap.add_argument("--id", action="append", help="only these endpoint ids") + ap.add_argument("--write", action="store_true", help="write examples/.json for passes") + ap.add_argument("--extended", action="store_true", + help="verify .extended.yaml instead of the curated .yaml " + "(only its entries that carry a test_request are callable)") + ap.add_argument("--via-treg", metavar="TOOL", + help="route through a treg registry's /call proxy (tool name) using ~/.treg/config.json " + "instead of a raw credential — for OAuth providers whose token treg holds") + args = ap.parse_args() + + cred = os.environ.get("TREG_CATALOG_CRED") + if not cred and not args.via_treg: + print("TREG_CATALOG_CRED not set (or use --via-treg )", file=sys.stderr) + return 2 + + sys.path.insert(0, str(ROOT / "src")) + from treg.oauth_providers import REGISTRY + + prov = REGISTRY.get(args.service) + if prov is None: + print(f"unknown provider '{args.service}'", file=sys.stderr) + return 2 + suffix = ".extended.yaml" if args.extended else ".yaml" + data = yaml.safe_load((CATALOG / f"{args.service}{suffix}").read_text()) + + headers, base_query, call_base = {}, {}, None + if args.via_treg: + cfg = json.loads((Path.home() / ".treg" / "config.json").read_text()) + call_base = cfg["base_url"].rstrip("/") + f"/call/{args.via_treg}" + headers["X-Treg-Token"] = cfg["token"] # treg's own auth header, not Bearer + if cfg.get("active_org"): + headers["X-Treg-Org"] = str(cfg["active_org"]) + else: + secret = cred + if prov.token_encode == "base64" and ":" in secret: + # curator may paste the raw login:password; the registry stores it base64ed + secret = base64.b64encode(secret.encode()).decode() + if prov.token_location == "query": + base_query[prov.token_param] = prov.token_format.format(secret=secret) + else: + headers[prov.token_header or "Authorization"] = (prov.token_format or "Bearer {secret}").format(secret=secret) + + failures = 0 + today = date.today().isoformat() + with httpx.Client(timeout=60, follow_redirects=True) as client: + for ep in data.get("endpoints", []): + eid = ep["id"] + if args.id and eid not in args.id: + continue + treq = ep.get("test_request") + if not treq: + continue # extended entries without one were never callable; nothing to replay + path = ep["path"] + for k, v in (treq.get("pathParams") or {}).items(): + path = path.replace("{%s}" % k, str(v)) + base = (call_base or prov.base_url.rstrip("/")) + # the extended tier stores each route exactly as the provider's spec spells it, which + # for DataForSEO includes the /v3 that base_url already ends with — joining blindly + # would request /v3/v3/… + prefix = urlsplit(base).path.rstrip("/") + if prefix and path.startswith(prefix + "/"): + path = path[len(prefix):] + url = base + path + params = {**base_query, **(treq.get("queryParams") or {})} + body = treq.get("body") + form = body if treq.get("bodyType") == "form" else None + if form is not None and base_query: + # form-bodied routes read the credential from the body; leaving it in the query + # as well makes the provider reject the call (verified live on Just One API) + form = {**base_query, **form} + params = {k: v for k, v in params.items() if k not in base_query} + try: + resp = client.request(ep["method"], url, params=params, headers=headers, + data=form, + json=body if body is not None and form is None else None) + except httpx.HTTPError as exc: + print(f"FAIL {eid} — transport error: {exc}") + failures += 1 + continue + + ok = 200 <= resp.status_code < 300 + detail = f"http {resp.status_code}" + doc = None + try: + doc = resp.json() + except ValueError: + pass + expect = ep.get("expect") + if ok and expect and doc is not None: + got = dig(doc, expect["json_path"]) + ok = got == expect.get("equals") + detail += f", {expect['json_path']}={got!r}" + if ok: + note = "" + if args.write and doc is not None: + ex_rel = ep.get("example_response") or f"examples/{eid}.json" + out = CATALOG / ex_rel + out.parent.mkdir(parents=True, exist_ok=True) + # shrink recursively until under the cap — never byte-slice serialized JSON + # (a mid-token slice writes an unparseable example file) + shrunk, max_str = truncate(doc), MAX_STR + dump = json.dumps(shrunk, indent=2, ensure_ascii=False) + while len(dump.encode()) > MAX_BYTES and max_str > 15: + max_str //= 4 + shrunk = _shrink_strings(shrunk, max_str) + dump = json.dumps(shrunk, indent=2, ensure_ascii=False) + out.write_text(dump + "\n") + note = f" -> {ex_rel}" + print(f"PASS {eid} ({detail}){note} [stamp verified: {today}]") + else: + snippet = (resp.text or "")[:160].replace("\n", " ") + print(f"FAIL {eid} — {detail}: {snippet}") + failures += 1 + return 1 if failures else 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/catalog_verify_extended.py b/scripts/catalog_verify_extended.py new file mode 100644 index 00000000..de0568f4 --- /dev/null +++ b/scripts/catalog_verify_extended.py @@ -0,0 +1,675 @@ +#!/usr/bin/env python3 +"""Live-verify the EXTENDED tier in bulk, under a hard spend cap. + + TREG_CATALOG_CRED='' uv run python scripts/catalog_verify_extended.py tikhub \ + --budget 1.80 [--platform tiktok] [--limit 50] [--dry-run] [--refresh] + +scripts/catalog_verify.py is the CORE-tier tool: a handful of hand-written endpoints, one line of +output each, a human reading every result. That does not survive contact with 1385 machine- +generated entries, where the questions are different — what does this cost, what can I afford to +skip, and how big does the repo get. This script is the extended-tier answer to those: + + * **Spend cap, enforced before the call.** Entries run CHEAPEST FIRST and the running total is + the sum of the per-endpoint `cost` of the calls that SUCCEEDED (TikHub bills per_success, so a + failure is free). The next call is only made if its price still fits in `--budget`. Cheapest- + first is not just safety: it buys the most verified endpoints per dollar, and it means a run + that stops early stops on the expensive tail, not somewhere arbitrary. + * **The account's own balance is the ground truth.** The provider's balance is read before and + after and printed as the real spend, because a rate card can be stale in a way our arithmetic + cannot detect. Where the provider states a per-call charge in its response (DataForSEO's + `tasks.0.cost`), that number is also written onto the endpoint as `observed_cost` — what it + ACTUALLY cost, next to the `cost` rate card that says what it should. Balance arithmetic cannot + separate this run's spend from anything else touching the same key, so a summed observed_cost + is the defensible figure for a report. + * **Examples trimmed far harder than core.** 1385 files at core's 10 KB cap would add ~14 MB of + JSON to the repo. Here: arrays keep ONE item, strings clip at 200 chars, 2 KB per document. + An extended example exists to show the response SHAPE; the core tier is where a full-fidelity + example belongs. + * **Deterministic and re-runnable.** Results are written back into the yaml — `verified` + a + `example_response` path on a pass, a one-line `unverified: http …` on a failure. A + re-run skips whatever already carries `verified` (use `--refresh` to re-test those too), so + interrupting a run and restarting it resumes rather than re-billing. + +Two retry rules, both free because only a successful call is billed: + + * a failure that reads like a flaky upstream ("please retry", a timeout, any 5xx, a 429, a + transport error) is retried up to --retry-attempts times. This is not politeness, it is the + single highest-value lever the script has: TikHub's LinkedIn family went from 8% to 90% + verified across four passes with no change other than being asked again. Conversion per pass + decays fast (672, +30, +16, +14, +2 on tikhub) — a pass that converts ~2 is convergence, and + what is left is genuinely broken rather than flaky. + * a 4xx on an endpoint whose page-size knob we clamped is retried once with the documented + values, which separates "our test request was wrong" from "this endpoint is broken". If that + retry passes, the working `test_request` is what gets written back. +""" + +from __future__ import annotations + +import argparse +import base64 +import json +import os +import re +import sys +import threading +import time +from concurrent.futures import ThreadPoolExecutor +from datetime import date +from pathlib import Path +from urllib.parse import urlsplit + +import httpx +import yaml + +ROOT = Path(__file__).resolve().parent.parent +CATALOG = ROOT / "src" / "treg" / "catalog" + +# far tighter than catalog_verify.py's core-tier caps — see the module docstring +MAX_STR = 200 +MAX_LIST = 1 +MAX_BYTES = 2_000 + +RATE_LIMIT_PER_SEC = 8.0 # provider allows 10/s; leave headroom so we are never the cause of a 429 + +# Per-call ceiling. Was 60s, which sat right on top of the slowest real endpoint we have measured: +# DataForSEO's /v3/merchant/amazon/products/live/advanced answered in 55.3s on a live probe. A +# timeout gets recorded as the endpoint's verdict, so a margin that thin turns a slow-but-working +# route into an intermittent "unverified" — the failure mode this file exists to avoid. +CALL_TIMEOUT = 180 + + +class Throttle: + """Global request pacer shared by the worker threads.""" + + def __init__(self, per_sec: float) -> None: + self._gap = 1.0 / per_sec + self._lock = threading.Lock() + self._next = 0.0 + + def wait(self) -> None: + with self._lock: + now = time.monotonic() + at = max(now, self._next) + self._next = at + self._gap + time.sleep(max(0.0, at - now)) + + +MIN_STR = 60 # a string clipped below this stops showing what the field CONTAINS + + +def trim(node, str_max: int = MAX_STR, depth_max: int = 99, depth: int = 0, keys_max: int = 99): + if isinstance(node, dict): + if depth >= depth_max: + return f"… {{{len(node)} field(s) truncated}}" + out = {k: trim(v, str_max, depth_max, depth + 1, keys_max) + for k, v in list(node.items())[:keys_max]} + if len(node) > keys_max: + out["…"] = f"{len(node) - keys_max} more field(s) truncated" + return out + if isinstance(node, list): + if depth >= depth_max: + return f"… [{len(node)} item(s) truncated]" + out = [trim(v, str_max, depth_max, depth + 1, keys_max) for v in node[:MAX_LIST]] + if len(node) > MAX_LIST: + out.append(f"… {len(node) - MAX_LIST} more item(s) truncated") + return out + if isinstance(node, str) and len(node) > str_max: + return node[:str_max] + f"… [{len(node)} chars]" + return node + + +def render_example(doc) -> str: + """Trim `doc` until it serialises under MAX_BYTES. Never byte-slices the JSON. + + Depth is given up before string length. What an extended example is FOR is showing the shape + of the response — which fields exist and what a value looks like — so clipping every string to + a dozen characters to preserve six levels of nesting gets the trade backwards: it keeps the + skeleton and throws away the only part that tells you what the endpoint returns. + """ + for str_max in (MAX_STR, 120, MIN_STR): + for depth_max, keys_max in ((99, 99), (6, 40), (5, 25), (4, 15), (3, 10)): + dump = json.dumps(trim(doc, str_max, depth_max, 0, keys_max), indent=2, ensure_ascii=False) + if len(dump.encode()) <= MAX_BYTES: + return dump + "\n" + return dump + "\n" + + +def cost_of(ep: dict) -> float: + """What the next call to this endpoint should be budgeted at. + + The rate card (`cost`) is preferred, but several providers publish prices per API family rather + than per route and their entries carry none at all — every DataForSEO entry, for one. A missing + price silently reads as FREE, which makes `--budget` inert: a run can queue the whole file at an + estimated $0.000 and still spend real money, with only the balance readback noticing afterwards. + So fall back to `observed_cost`, what the provider charged last time. It is the better number + anyway — measured rather than transcribed — and it makes the cap bind from the second run on. + """ + c = ep.get("cost") or {} + try: + listed = float(c.get("value") or 0.0) + except (TypeError, ValueError): + listed = 0.0 + if listed: + return listed + observed = ep.get("observed_cost") + return float(observed) if isinstance(observed, (int, float)) else 0.0 + + +def auth_headers(service: str, cred: str) -> tuple[dict, dict, str]: + sys.path.insert(0, str(ROOT / "src")) + from treg.oauth_providers import REGISTRY + + prov = REGISTRY.get(service) + if prov is None: + raise SystemExit(f"unknown provider '{service}'") + headers, query = {}, {} + secret = cred + if prov.token_encode == "base64" and ":" in secret: + secret = base64.b64encode(secret.encode()).decode() + if prov.token_location == "query": + query[prov.token_param] = prov.token_format.format(secret=secret) + else: + headers[prov.token_header or "Authorization"] = ( + prov.token_format or "Bearer {secret}").format(secret=secret) + return headers, query, prov.base_url.rstrip("/") + + +def join(base: str, path: str) -> str: + """base_url + a catalog path, without doubling a shared prefix. + + The extended tier stores each route exactly as the provider's spec spells it, and DataForSEO's + include the `/v3` that its `base_url` already ends with — so a blind join requests + `/v3/v3/serp/...` and every call 404s. Same fix as catalog_verify.py; a no-op for providers + whose base_url is a bare host. + """ + prefix = urlsplit(base).path.rstrip("/") + if prefix and path.startswith(prefix + "/"): + path = path[len(prefix):] + return base + path + + +# service -> (free route that reports the account balance, dotted path to the number) +BALANCE_ROUTE = { + "tikhub": ("/api/v1/tikhub/user/get_user_info", "user_data.balance"), + "dataforseo": ("/v3/appendix/user_data", "tasks.0.result.0.money.balance"), +} + + +def dig(doc, dotted: str): + """Walk a dotted path, where a numeric segment indexes a list (`tasks.0.result.0.money`).""" + cur = doc + for part in dotted.split("."): + if isinstance(cur, dict): + cur = cur.get(part) + elif isinstance(cur, list) and part.lstrip("-").isdigit(): + try: + cur = cur[int(part)] + except IndexError: + return None + else: + return None + return cur + + +BALANCE_ATTEMPTS = 3 + + +def read_balance(client: httpx.Client, service: str, base: str, headers: dict, query: dict): + """The account's own balance, or None if the provider won't say. + + Retried, because a balance endpoint under repeated polling is exactly where providers get + flaky: DataForSEO answers with a null `result` when polled in quick succession. A None here is + not fatal — the run's own arithmetic still bounds the spend — but it silently removes the only + independent check on that arithmetic, so it is worth asking more than once. + """ + route = BALANCE_ROUTE.get(service) + if not route: + return None + for attempt in range(BALANCE_ATTEMPTS): + try: + r = client.get(join(base, route[0]), headers=headers, params=query, timeout=30) + value = dig(r.json(), route[1]) + if isinstance(value, (int, float)): + return value + except Exception: + pass + time.sleep(1.0 * (attempt + 1)) + return None + + +def unclamped(ep: dict) -> dict | None: + """`test_request` with page-size knobs restored to their documented values, or None. + + Used for the single free retry after a 4xx: it tells apart an endpoint that rejects our small + page size from an endpoint that is actually broken. + """ + test = json.loads(json.dumps(ep.get("test_request") or {})) + schema = ((ep.get("input") or {}).get("queryParams") or {}) + changed = False + for name, value in list((test.get("queryParams") or {}).items()): + doc_example = (schema.get(name) or {}).get("example") + if doc_example not in (None, "") and doc_example != value: + test["queryParams"][name] = doc_example + changed = True + return test if changed else None + + +# TikHub answers a failed SCRAPE with 400 and this text — the request was well-formed, their +# upstream fetch just did not come back. LinkedIn measured ~1-in-3 on first attempt and passed on +# a retry, so a single attempt would have written off 44 working endpoints as broken. Retrying is +# free (nothing is billed unless the call succeeds), which makes "try again" strictly better than +# recording a verdict we know is unreliable. +RETRYABLE = re.compile(r"please retry|timed? ?out|temporarily|try again|too fast", re.I) +RETRY_ATTEMPTS = 3 +# TikHub rate-limits PER ROUTE at 1 request/second, independently of the account-wide 10/s that +# RATE_LIMIT_PER_SEC paces. The global throttle cannot help here — it spaces consecutive requests +# across DIFFERENT routes — so any retry, which by definition hits the same route again, must wait +# out that second itself. A sub-second backoff turns a retry into a guaranteed 429, and the 429 +# then gets recorded as if the endpoint had failed. That mistake cost 66 endpoints a real verdict +# on the first full run. +PER_ROUTE_GAP = 1.3 + + +def retryable(detail: str, status: int | None) -> bool: + if status is not None and (status >= 500 or status == 429): + return True + if detail.startswith("transport"): + return True + return bool(RETRYABLE.search(detail)) + + +def call_with_retries(client: httpx.Client, base: str, ep: dict, test: dict, headers: dict, + query: dict, throttle: Throttle, max_attempts: int = RETRY_ATTEMPTS + ) -> tuple[bool, str, object, int]: + """`call`, repeated while the failure looks like a flaky upstream rather than a verdict.""" + attempts = 0 + for attempt in range(max_attempts): + attempts += 1 + ok, detail, doc, status = call(client, base, ep, test, headers, query, throttle) + if ok or not retryable(detail, status): + return ok, detail, doc, attempts + time.sleep(PER_ROUTE_GAP * (attempt + 1)) + return ok, detail, doc, attempts + + +def call(client: httpx.Client, base: str, ep: dict, test: dict, headers: dict, query: dict, + throttle: Throttle) -> tuple[bool, str, object, int | None]: + path = ep["path"] + for k, v in (test.get("pathParams") or {}).items(): + path = path.replace("{%s}" % k, str(v)) + params = {**query, **(test.get("queryParams") or {})} + body = test.get("body") + form = body if test.get("bodyType") == "form" else None + if form is not None and query: + # these routes take the token in the form body; leaving a copy in the query string makes + # the provider reject the call with a misleading auth error (verified live on Just One API) + form = {**query, **form} + params = {k: v for k, v in params.items() if k not in query} + throttle.wait() + started = time.monotonic() + try: + resp = client.request(ep["method"], join(base, path), params=params, headers=headers, + data=form, json=body if body is not None and form is None else None, + timeout=CALL_TIMEOUT) + except httpx.HTTPError as exc: + return False, f"transport {type(exc).__name__}", None, None + finally: + LAST_ELAPSED.seconds = round(time.monotonic() - started, 2) + doc = None + try: + doc = resp.json() + except ValueError: + pass + ok = 200 <= resp.status_code < 300 + # TikHub answers 200 for the HTTP layer and puts the real verdict in the body's `code` + if ok and isinstance(doc, dict) and isinstance(doc.get("code"), int) and doc["code"] not in (0, 200): + ok, detail = False, f"http {resp.status_code} body-code {doc['code']}" + else: + detail = f"http {resp.status_code}" + if not ok: + if doc is None: + msg = resp.text or "" + else: + msg = (doc.get("detail") if isinstance(doc, dict) else None) or doc + if isinstance(msg, dict): + msg = msg.get("message") or msg.get("detail") or "" + snippet = re.sub(r"\s+", " ", str(msg)).strip()[:80] + if snippet: + detail = f"{detail}: {snippet}" + return ok, detail, doc, resp.status_code + + +def save(path: Path, data: dict, header: str) -> None: + body = yaml.safe_dump(data, sort_keys=False, allow_unicode=True, width=4096, + default_flow_style=False) + path.write_text(header + body) + + +RESULT_FIELDS = ("verified", "unverified", "skipped", "example_response", "test_request", + "observed_cost", "observed_time") + + +class _Elapsed(threading.local): + """Per-thread stopwatch for the last request, so `call` keeps its return signature.""" + seconds = None + + +LAST_ELAPSED = _Elapsed() +# `observed_time` is OUR wall-clock around the request, not a number read out of the response. +# Only DataForSEO reports its own duration, and TikHub's `time` field is a TIMESTAMP +# ("2026-07-27 23:27:48") — an extractor that trusted the field name would write a date string +# into a numeric column. Client-side is also the number that matters here: CALL_TIMEOUT applies to +# our client, so what tells us a route is near the ceiling is how long WE waited. Recorded because +# a timeout is written down as the endpoint's verdict, and DataForSEO's +# /v3/merchant/amazon/products/live/advanced answered in 55.3s and 26.2s on two identical calls — +# a 2x spread that would have made a healthy route intermittently "unverified" under the old 60s. + +# service -> how to read what the provider says THIS call cost, from its own response. +# `cost` in the catalog is a rate card: what the price list claims, entered by hand and liable to +# age. `observed_cost` is what the provider actually charged, straight from the response that +# produced the example — so the two can be compared, a run's spend can be summed without balance +# arithmetic (which cannot separate our calls from anyone else's on the same key), and the router +# later has a real per-endpoint price to rank on. +# TikHub is deliberately absent: it publishes no per-call charge in the response body, so its +# entries carry no observed_cost and its rate card stays the only figure available. +OBSERVED_COST = { + "dataforseo": lambda doc: dig(doc, "tasks.0.cost"), +} + + +def retrim(service: str) -> int: + """Re-apply the current size caps to already-captured examples, without calling anything. + + The captured JSON is on disk, so shrinking it further is a local transformation — re-running + the endpoint to get a smaller file would pay the provider a second time for a response we + already have. Only genuinely unreadable files are worth a re-capture, and this reports those + rather than deleting them. + + Trimming is one-way: this cannot restore detail an earlier, coarser pass already dropped. It + exists to bring old captures under a cap that got stricter, which is the only direction that + matters for repo size. + """ + file = CATALOG / f"{service}.extended.yaml" + eps = (yaml.safe_load(file.read_text()) or {}).get("endpoints") or [] + shrunk = saved = unreadable = 0 + broken: list[str] = [] + for ep in eps: + rel = ep.get("example_response") + if not rel: + continue + path = CATALOG / rel + if not path.is_file(): + continue + before = path.stat().st_size + if before <= MAX_BYTES: + continue + try: + doc = json.loads(path.read_text()) + except ValueError: + unreadable += 1 + broken.append(ep["id"]) + continue + path.write_text(render_example(doc)) + after = path.stat().st_size + if after < before: + shrunk += 1 + saved += before - after + print(f"re-trimmed {shrunk} example(s), reclaimed {saved / 1024:.0f} KB, $0 spent") + if unreadable: + print(f"{unreadable} unreadable file(s) need a re-capture: {broken[:5]}", file=sys.stderr) + over = [ep["id"] for ep in eps if ep.get("example_response") + and (CATALOG / ep["example_response"]).is_file() + and (CATALOG / ep["example_response"]).stat().st_size > MAX_BYTES] + print(f"{len(over)} example(s) still over {MAX_BYTES} B" + (f", e.g. {over[:3]}" if over else "")) + return 0 + + +def merge_results(service: str, paths: list[Path]) -> int: + """Fold batch result files into `.extended.yaml`. The ONLY writer of that file. + + Splitting a provider across concurrent agents means splitting the endpoints, not the file: two + processes rewriting one yaml lose each other's work, and the loss is silent because both + produce a valid document. So a batch run takes `--out` and writes only its own results, and one + merge folds them in. Example RESPONSES need no such care — each is its own file named after the + endpoint that produced it, so disjoint batches never touch the same one. + + A batch may only report on endpoints it was given; an id that is not already in the yaml is + refused rather than added, because the endpoint list is catalog_ingest.py's to own. + """ + file = CATALOG / f"{service}.extended.yaml" + text = file.read_text() + header = text[: text.index("\nprovider:") + 1] if "\nprovider:" in text else "" + data = yaml.safe_load(text) + by_id = {e["id"]: e for e in data["endpoints"]} + + applied, unknown, conflicts = 0, [], [] + seen: dict[str, Path] = {} + for path in paths: + payload = json.loads(path.read_text()) + if payload.get("service") != service: + raise SystemExit(f"{path}: results are for '{payload.get('service')}', not '{service}'") + for res in payload.get("results") or []: + eid = res.get("id") + ep = by_id.get(eid) + if ep is None: + unknown.append(eid) + continue + if eid in seen: + conflicts.append(f"{eid} (in {seen[eid].name} and {path.name})") + continue + seen[eid] = path + for field in RESULT_FIELDS: + if field in res: + # an EMPTY value counts as absent, never as a claim. A state key present but + # blank passes a "has a state" check while saying nothing, which is worse than + # no key at all — the entry looks handled and isn't. (Same bug the dataforseo + # run hit when a re-run dropped a curated reason string.) + if res[field] is None or res[field] == "": + ep.pop(field, None) + else: + ep[field] = res[field] + if res.get("verified"): + ep.pop("untestable", None) + # exactly one outcome state survives a merge: a batch that reports an endpoint as + # verified must not leave last run's `unverified` sitting beside it + states = {"verified", "unverified", "skipped"} + if states & set(res): + for stale in states - {s for s in res if res.get(s)}: + ep.pop(stale, None) + applied += 1 + + if conflicts: + # two batches claiming one endpoint means the split was wrong; merging either result would + # hide that, so refuse and let the caller fix the partition + raise SystemExit("overlapping batches: " + "; ".join(conflicts[:10])) + if unknown: + print(f"warning: {len(unknown)} unknown id(s) ignored, e.g. {unknown[:3]}", file=sys.stderr) + save(file, data, header) + print(f"merged {applied} result(s) from {len(paths)} batch file(s) into {file.relative_to(ROOT)}") + print(f"{sum(1 for e in data['endpoints'] if e.get('verified'))} of {len(data['endpoints'])} verified") + return 0 + + +def main(argv: list[str]) -> int: + ap = argparse.ArgumentParser(description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("service") + ap.add_argument("--budget", type=float, default=1.0, help="max USD to spend this run") + ap.add_argument("--platform", action="append", help="only these platforms") + ap.add_argument("--limit", type=int, help="stop after N calls") + ap.add_argument("--max-cost", type=float, help="skip endpoints priced above this") + ap.add_argument("--concurrency", type=int, default=5) + ap.add_argument("--refresh", action="store_true", help="re-test endpoints already verified") + ap.add_argument("--dry-run", action="store_true", help="plan and cost only; make no calls") + ap.add_argument("--out", type=Path, metavar="FILE", + help="write results as JSON here instead of touching the yaml (batch mode)") + ap.add_argument("--merge", type=Path, nargs="+", metavar="FILE", + help="fold --out result files into the yaml; makes no calls") + ap.add_argument("--retry-attempts", type=int, default=RETRY_ATTEMPTS, + help="attempts per endpoint when the failure looks transient. Raising this is " + "free (only successes bill) and pays off on providers whose scraper is " + "flaky rather than broken — TikHub's Douyin/LinkedIn families convert " + "steadily across repeated attempts") + ap.add_argument("--retrim", action="store_true", + help="re-apply the size caps to captured examples locally; makes no calls") + args = ap.parse_args(argv) + + if args.merge: + return merge_results(args.service, args.merge) + if args.retrim: + return retrim(args.service) + + cred = os.environ.get("TREG_CATALOG_CRED") + if not cred and not args.dry_run: + print("TREG_CATALOG_CRED not set", file=sys.stderr) + return 2 + + file = CATALOG / f"{args.service}.extended.yaml" + text = file.read_text() + header = text[: text.index("\nprovider:") + 1] if "\nprovider:" in text else "" + data = yaml.safe_load(text) + eps = data["endpoints"] + + todo = [e for e in eps if e.get("test_request") is not None] + if args.platform: + todo = [e for e in todo if e["platform"] in set(args.platform)] + if not args.refresh: + todo = [e for e in todo if not e.get("verified")] + if args.max_cost is not None: + todo = [e for e in todo if cost_of(e) <= args.max_cost] + todo.sort(key=lambda e: (cost_of(e), e["platform"], e["path"])) + + print(f"{len(eps)} endpoint(s); {sum(1 for e in eps if e.get('test_request') is not None)} testable; " + f"{len(todo)} queued this run") + print(f"queued at list price: ${sum(cost_of(e) for e in todo):.3f} budget ${args.budget:.2f}") + if args.dry_run: + return 0 + + headers, query, base = auth_headers(args.service, cred) + spent = 0.0 + stats = {"pass": 0, "fail": 0, "skipped_budget": 0, "retried": 0} + touched: set[str] = set() + lock = threading.Lock() + throttle = Throttle(RATE_LIMIT_PER_SEC) + stop = threading.Event() + calls = 0 + + with httpx.Client(timeout=60, follow_redirects=True) as client: + before = read_balance(client, args.service, base, headers, query) + print(f"balance before: {before}") + + def work(ep: dict) -> None: + nonlocal spent, calls + price = cost_of(ep) + with lock: + if stop.is_set(): + return + if spent + price > args.budget: + # the FOURTH state (convention shared with the dataforseo/justoneapi run): a + # usable test request exists and the call was simply never made. Calling this + # `unverified` would invent a failure that never happened, and `untestable` would + # be a false claim about the endpoint. A `skipped` entry clears with money and a + # re-run; the other two need a human. + stats["skipped_budget"] += 1 + if ep.get("verified"): + # already proven by an earlier run and merely not RE-checked by this one. + # Downgrading it to `skipped` would throw away a real verification, and + # emitting both states puts a contradiction in the --out file that a later + # --merge would fold into the yaml. Leave it alone and report nothing. + return + ep["skipped"] = (f"not called — ${price:.4f} would exceed this run's " + f"${args.budget:.2f} budget") + ep.pop("unverified", None) + touched.add(ep["id"]) + return + if args.limit and calls >= args.limit: + stop.set() + return + calls += 1 + # reserve the price up front so N concurrent workers cannot jointly overshoot; + # it is released below if the call fails (failures are not billed) + spent += price + ok, detail, doc, tries = call_with_retries(client, base, ep, ep["test_request"], + headers, query, throttle, + args.retry_attempts) + if not ok and detail.startswith("http 4"): + alt = unclamped(ep) + if alt: + time.sleep(PER_ROUTE_GAP) # same route again — clear its 1/s window first + stats["retried"] += 1 + ok2, detail2, doc2, _ = call_with_retries(client, base, ep, alt, headers, + query, throttle, args.retry_attempts) + if ok2: + ep["test_request"] = alt + ok, detail, doc = ok2, detail2, doc2 + with lock: + if not ok: + spent -= price + if ok and doc is not None: + rel = f"examples/{ep['id']}.json" + out = CATALOG / rel + out.parent.mkdir(parents=True, exist_ok=True) + out.write_text(render_example(doc)) + ep["verified"] = date.today().isoformat() + ep["example_response"] = rel + observed = (OBSERVED_COST.get(args.service) or (lambda _d: None))(doc) + if isinstance(observed, (int, float)): + ep["observed_cost"] = observed + if isinstance(LAST_ELAPSED.seconds, (int, float)): + ep["observed_time"] = LAST_ELAPSED.seconds + ep.pop("unverified", None) + ep.pop("skipped", None) + with lock: + touched.add(ep["id"]) + stats["pass"] += 1 + else: + # never write a blank state: an empty reason satisfies "has a state" while telling + # the next reader nothing, so fall back to the bare status rather than to silence + ep["unverified"] = detail.strip() or "failed with no status or message" + ep.pop("verified", None) + ep.pop("example_response", None) + ep.pop("skipped", None) + with lock: + touched.add(ep["id"]) + stats["fail"] += 1 + done = stats["pass"] + stats["fail"] + if done % 50 == 0: + print(f" … {done} called, {stats['pass']} pass, ${spent:.3f} spent", flush=True) + + try: + with ThreadPoolExecutor(max_workers=args.concurrency) as pool: + list(pool.map(work, todo)) + except KeyboardInterrupt: + stop.set() + print("interrupted — writing partial results", file=sys.stderr) + + after = read_balance(client, args.service, base, headers, query) + + if args.out: + # batch mode: report only what THIS run called, so the merge cannot resurrect a stale + # verdict for an endpoint another batch owns + results = [{"id": e["id"], **{f: e.get(f) for f in RESULT_FIELDS}} + for e in eps if e["id"] in touched] + args.out.parent.mkdir(parents=True, exist_ok=True) + args.out.write_text(json.dumps( + {"service": args.service, "generated": date.today().isoformat(), + "platforms": sorted({e["platform"] for e in todo}), "results": results}, + indent=2, ensure_ascii=False) + "\n") + else: + save(file, data, header) + print(f"\nPASS {stats['pass']} FAIL {stats['fail']} skipped-for-budget {stats['skipped_budget']}" + f" retried {stats['retried']}") + print(f"estimated spend ${spent:.3f}; balance {before} -> {after}" + + (f" (actual ${before - after:.3f})" if isinstance(before, (int, float)) + and isinstance(after, (int, float)) else "")) + if args.out: + print(f"wrote {len(results)} result(s) to {args.out}" + f" — {file.name} NOT touched; apply with --merge {args.out}") + else: + print(f"{sum(1 for e in eps if e.get('verified'))} of {len(eps)} endpoint(s) now verified;" + f" wrote {file.relative_to(ROOT)}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main(sys.argv[1:]))