mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
582 lines
27 KiB
Python
582 lines
27 KiB
Python
"""Shared test fixtures.
|
|
|
|
The "upstream" is a tiny in-process ASGI echo app. The registry's shared httpx client is
|
|
pointed at it via ASGITransport, so the relay path runs for real, just without a socket.
|
|
The `clients` fixture also registers a user and authes the client by default.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import atexit
|
|
import glob
|
|
import os
|
|
import socket
|
|
import tempfile
|
|
|
|
# Tests and their CLI subprocesses must never emit production analytics.
|
|
os.environ["TREG_TELEMETRY"] = "0"
|
|
|
|
# Isolate the test DB from any .env / running dev server BEFORE importing treg (the engine is
|
|
# built at import time). A real env var overrides the .env file in pydantic-settings.
|
|
# TREG_TEST_DB_URL (not TREG_DATABASE_URL - a stray production URL in a shell must never become the
|
|
# test target) points the suite at a specific database, e.g. the Postgres CI job.
|
|
# Otherwise every test process gets its OWN sqlite file, named by pid: processes sharing one file
|
|
# wipe each other's rows mid-test, whether they are xdist workers of one run (1,022 errors on the
|
|
# first parallel run) or two runs started side by side in one checkout.
|
|
_worker = os.environ.get("PYTEST_XDIST_WORKER", "")
|
|
# The files live under the system temp dir, NOT the repo root: sixteen 600 KB databases rewritten
|
|
# on every run kept editors' file watchers busy re-indexing the working tree.
|
|
_db_dir = os.path.join(tempfile.gettempdir(), "treg-tests")
|
|
os.makedirs(_db_dir, exist_ok=True)
|
|
_db_file = os.path.join(_db_dir, f"treg-test-{os.getpid()}.db")
|
|
atexit.register(lambda: [os.remove(p) for p in glob.glob(_db_file + "*")])
|
|
_default = f"sqlite+aiosqlite:///{_db_file}"
|
|
|
|
|
|
def _per_worker_postgres(url: str) -> str:
|
|
"""An xdist worker against Postgres gets its own database (treg_test_gw0, ...), created on first
|
|
use: `reset_db()` empties every table, so workers sharing one database wipe each other's rows."""
|
|
import asyncio
|
|
import asyncpg
|
|
from sqlalchemy.engine import make_url
|
|
|
|
base = make_url(url)
|
|
name = f"{base.database}_{_worker}"
|
|
|
|
async def ensure() -> None:
|
|
dsn = base.set(drivername="postgresql").render_as_string(hide_password=False)
|
|
conn = await asyncpg.connect(dsn)
|
|
try:
|
|
if not await conn.fetchval("SELECT 1 FROM pg_database WHERE datname = $1", name):
|
|
await conn.execute(f'CREATE DATABASE "{name}"')
|
|
finally:
|
|
await conn.close()
|
|
|
|
asyncio.run(ensure())
|
|
return base.set(database=name).render_as_string(hide_password=False)
|
|
|
|
|
|
_test_db_url = os.environ.get("TREG_TEST_DB_URL")
|
|
if _test_db_url and _worker and _test_db_url.startswith("postgresql"):
|
|
# Resolved once per process: tests that import `tests.conftest` load this module a second time.
|
|
_test_db_url = os.environ.get("TREG_TEST_DB_URL_WORKER") or _per_worker_postgres(_test_db_url)
|
|
os.environ["TREG_TEST_DB_URL_WORKER"] = _test_db_url
|
|
os.environ["TREG_DATABASE_URL"] = _test_db_url or _default
|
|
# Replica tests opt in explicitly; never inherit a real replica from the shell or .env.
|
|
os.environ["TREG_READ_DATABASE_URL"] = ""
|
|
os.environ["TREG_EMAIL_DEV_MODE"] = "true" # tests need the returned OTP code (prod default is now False)
|
|
os.environ["TREG_RESEND_API_KEY"] = "" # never fire a real Resend send from the test suite (send_otp/send_invite skip when empty)
|
|
os.environ["TREG_RUN_ALLOWED_BINS"] = "sh,echo,true,false,cat,sleep,treg-nonexistent-bin-xyz" # allow the test CLIs for --server run tests
|
|
os.environ["TREG_PROXY_SSRF_CHECK"] = "false"
|
|
# The production default is OFF. The established connector tests and committed route snapshots test
|
|
# the enabled product surface; dedicated tests below also prove the disabled deployment shape.
|
|
os.environ["TREG_CLAUDE_CONNECTOR_ENABLED"] = "true"
|
|
# Blank every registry credential so the suite NEVER inherits a developer's real .env. Settings
|
|
# reads .env, and a real env var beats it — so without this, a machine with Google/X/LinkedIn
|
|
# credentials configured runs a different suite than CI, and provider tests pass or fail depending
|
|
# on whose laptop they're on (google-ads autoprovisions only when a developer token is present).
|
|
# Tests that need a credential set it explicitly via monkeypatch.
|
|
for _k in (
|
|
"GOOGLE_CLIENT_ID", "GOOGLE_CLIENT_SECRET", "GOOGLE_ADS_DEVELOPER_TOKEN",
|
|
"LINKEDIN_CLIENT_ID", "LINKEDIN_CLIENT_SECRET",
|
|
"X_CLIENT_ID", "X_CLIENT_SECRET", "SLACK_CLIENT_ID", "SLACK_CLIENT_SECRET",
|
|
"TIKTOK_CLIENT_KEY", "TIKTOK_CLIENT_SECRET",
|
|
"META_CLIENT_ID", "META_CLIENT_SECRET",
|
|
"POSTHOG_KEY", "ADS_CONV_REFRESH_TOKEN",
|
|
"INSTAGRAM_CLIENT_ID", "INSTAGRAM_CLIENT_SECRET",
|
|
# …and the tier-4 platform keys + their allow-list. A developer's .env carries real, FUNDED keys:
|
|
# without this a suite run on their laptop could resolve tier 4 and spend actual money on the
|
|
# in-process upstream's echo. Tests that exercise tier 4 set both halves via monkeypatch.
|
|
"PLATFORM_KEY_TRYKITT", "PLATFORM_PROVIDERS", "PLATFORM_KEY_TIKHUB", "PLATFORM_KEY_DATAFORSEO", "PLATFORM_KEY_SCRAPECREATORS",
|
|
"PLATFORM_KEY_QUICKENRICH", "PLATFORM_KEY_PROSPEO", "PLATFORM_KEY_AIARK", "PLATFORM_KEY_WIZA",
|
|
"PLATFORM_KEY_SCRUBBY", "PLATFORM_KEY_ZEROBOUNCE", "PLATFORM_KEY_SUMBLE", "PLATFORM_KEY_MOLTSETS",
|
|
"PLATFORM_KEY_OPENMART", "PLATFORM_KEY_LIMADATA", "PLATFORM_KEY_HARVESTAPI", "PLATFORM_KEY_FETCHINIO",
|
|
"PLATFORM_KEY_DROPLEADS",
|
|
"PLATFORM_KEY_FINANCIALDATASETS",
|
|
"PLATFORM_KEY_ADYNTEL", "PLATFORM_EMAIL_ADYNTEL",
|
|
"PLATFORM_KEY_KEENABLE", "PLATFORM_KEY_OLOSTEP",
|
|
"PLATFORM_KEY_SEARCH1API",
|
|
):
|
|
os.environ[f"TREG_{_k}"] = "" # the test upstream is an in-process ASGI transport, not real DNS
|
|
|
|
import pytest # noqa: E402
|
|
from fastapi import FastAPI, Request # noqa: E402
|
|
from fastapi.responses import JSONResponse # noqa: E402
|
|
from httpx import ASGITransport, AsyncClient # noqa: E402
|
|
|
|
from treg import audit # noqa: E402
|
|
from treg.api import app # noqa: E402
|
|
from treg import archive # noqa: E402
|
|
from treg.config import get_settings # noqa: E402
|
|
from treg.infra.db import reset_db # noqa: E402
|
|
from treg.domain.identity import api_keys as managed_keys # noqa: E402
|
|
|
|
|
|
# The OTP-start + sandbox throttles (and the OTP codes) now live in the DB's `ephemeral` table, not in
|
|
# process-global dicts — so `reset_db()` (called by every client fixture) already clears them between
|
|
# tests. No separate rate-limit reset fixture is needed.
|
|
|
|
|
|
@pytest.fixture
|
|
def fake_getaddrinfo(monkeypatch):
|
|
"""Override named hosts only, leaving DB and other infrastructure DNS untouched.
|
|
|
|
An empty address list models an unresolvable host without querying external DNS.
|
|
"""
|
|
original = socket.getaddrinfo
|
|
|
|
def install(addresses: dict[str, list[str]]) -> None:
|
|
def resolve(host, port, *args, **kwargs):
|
|
if host not in addresses:
|
|
return original(host, port, *args, **kwargs)
|
|
if not addresses[host]:
|
|
raise socket.gaierror(socket.EAI_NONAME, "unresolvable")
|
|
return [
|
|
(socket.AF_INET6 if ":" in address else socket.AF_INET,
|
|
socket.SOCK_STREAM, 0, "",
|
|
(address, port or 0, 0, 0) if ":" in address else (address, port or 0))
|
|
for address in addresses[host]
|
|
]
|
|
monkeypatch.setattr(socket, "getaddrinfo", resolve)
|
|
|
|
return install
|
|
|
|
|
|
def make_upstream(hook_hits: list | None = None) -> FastAPI:
|
|
up = FastAPI()
|
|
|
|
@up.post("/token")
|
|
async def token() -> dict:
|
|
# stand-in OAuth token endpoint: serves both refresh + authorization_code exchanges.
|
|
return {"access_token": "REFRESHED", "refresh_token": "NEW-RT", "expires_in": 3600}
|
|
|
|
@up.get("/webmasters/v3/sites")
|
|
async def sites() -> dict:
|
|
# stand-in for a provider's resource-listing endpoint (GSC's shape), so connection
|
|
# discovery can be exercised without reaching Google.
|
|
return {
|
|
"siteEntry": [
|
|
{"siteUrl": "sc-domain:example.com", "displayName": "Example (production)"},
|
|
{"siteUrl": "https://staging.example/", "displayName": "Example (staging)"},
|
|
]
|
|
}
|
|
|
|
@up.post("/v25.0/oauth/access_token")
|
|
@up.get("/v25.0/oauth/access_token")
|
|
async def meta_token() -> dict:
|
|
# Meta's token endpoint, serving both the code exchange (POST) and the long-lived
|
|
# fb_exchange_token swap (GET). The ASGI transport routes every host here, so the real
|
|
# graph.facebook.com path must exist for a registry-mode Meta connect to complete.
|
|
return {"access_token": "META-TOKEN", "token_type": "bearer", "expires_in": 5183944}
|
|
|
|
@up.post("/oauth/access_token")
|
|
async def instagram_short_token() -> dict:
|
|
return {"access_token": "IG-SHORT-TOKEN", "user_id": 17841400000000000}
|
|
|
|
@up.get("/access_token")
|
|
async def instagram_long_token() -> dict:
|
|
return {"access_token": "IG-LONG-TOKEN", "token_type": "bearer", "expires_in": 5184000}
|
|
|
|
@up.get("/v25.0/me")
|
|
async def instagram_identity(request: Request):
|
|
if request.headers.get("authorization") != "Bearer IG-LONG-TOKEN":
|
|
return JSONResponse({"error": {"message": "invalid token"}}, status_code=401)
|
|
return {"user_id": "17841400000000000", "username": "direct_ig"}
|
|
|
|
@up.get("/me/accounts")
|
|
@up.get("/v25.0/me/accounts")
|
|
async def meta_pages(request: Request) -> dict:
|
|
# Meta's primary Page listing: what the user manages through a PERSONAL Page role. One row
|
|
# carries both the facebook shape (id/name) and the instagram shape (nested professional
|
|
# account), so both Meta providers can discover against the same stand-in.
|
|
page = {
|
|
"id": "PAGE-DIRECT", "name": "Directly Managed Page",
|
|
"instagram_business_account": {"id": "IG-DIRECT", "username": "direct_ig"},
|
|
}
|
|
if "access_token" in request.query_params.get("fields", ""):
|
|
page["access_token"] = "PAGE-TOKEN-DIRECT"
|
|
return {"data": [page]}
|
|
|
|
@up.get("/me/businesses")
|
|
@up.get("/v25.0/me/businesses")
|
|
async def meta_businesses(request: Request):
|
|
# Meta's Business walk (needs business_management): each business row nests owned_pages /
|
|
# client_pages whose entries are shaped like /me/accounts rows. PAGE-DIRECT reappears here
|
|
# (a personal-role Page is usually also Business-owned) to exercise dedup, and the
|
|
# agency-owned Page without a linked Instagram account must drop out of the IG picker.
|
|
# A token containing "noscope" emulates a connection that consented before
|
|
# business_management was in our scopes.
|
|
if "noscope" in request.headers.get("authorization", ""):
|
|
return JSONResponse(
|
|
{"error": {"message": "(#100) Missing Permission", "code": 100}}, status_code=400)
|
|
direct = {
|
|
"id": "PAGE-DIRECT", "name": "Directly Managed Page",
|
|
"instagram_business_account": {"id": "IG-DIRECT", "username": "direct_ig"},
|
|
}
|
|
client = {
|
|
"id": "PAGE-CLIENT", "name": "Agency Client Page",
|
|
"instagram_business_account": {"id": "IG-CLIENT", "username": "client_ig"},
|
|
}
|
|
if "access_token" in request.query_params.get("fields", ""):
|
|
direct["access_token"] = "PAGE-TOKEN-DIRECT"
|
|
client["access_token"] = "PAGE-TOKEN-CLIENT"
|
|
return {
|
|
"data": [
|
|
{"id": "BIZ-1", "owned_pages": {"data": [
|
|
direct,
|
|
{"id": "PAGE-NO-IG", "name": "Business Page Without Instagram"},
|
|
]}},
|
|
{"id": "BIZ-2", "client_pages": {"data": [
|
|
client,
|
|
]}},
|
|
]
|
|
}
|
|
|
|
@up.get("/v25.0/{page_id}/conversations")
|
|
async def instagram_conversations(page_id: str, request: Request):
|
|
if request.headers.get("authorization") == "Bearer IG-LONG-TOKEN":
|
|
return {
|
|
"auth": request.headers["authorization"],
|
|
"raw_path": request.scope.get("raw_path", b"").decode(),
|
|
"data": [{"id": "IG-DIRECT-CONVERSATION-1", "ig_user_id": page_id}],
|
|
}
|
|
# Messaging is the regression this Meta token split fixes: the user token is valid OAuth,
|
|
# but this edge accepts only the token of the Page linked to the selected Instagram account.
|
|
if request.headers.get("authorization") != "Bearer PAGE-TOKEN-DIRECT":
|
|
return JSONResponse(
|
|
{"error": {"message": "must be called with a Page access token", "code": 190}},
|
|
status_code=403,
|
|
)
|
|
if request.query_params.get("platform") != "instagram":
|
|
return {"data": []}
|
|
return {"data": [{"id": "IG-CONVERSATION-1", "page_id": page_id}]}
|
|
|
|
@up.post("/v25.0/{page_id}/subscribed_apps")
|
|
async def instagram_subscribe(page_id: str, request: Request):
|
|
expected_token = {
|
|
"PAGE-DIRECT": "PAGE-TOKEN-DIRECT",
|
|
"PAGE-CLIENT": "PAGE-TOKEN-CLIENT",
|
|
}.get(page_id)
|
|
if request.headers.get("authorization") != f"Bearer {expected_token}":
|
|
return JSONResponse(
|
|
{"error": {"message": "must be called with a Page access token", "code": 190}},
|
|
status_code=403,
|
|
)
|
|
subscribed_fields = request.query_params.get("subscribed_fields", "")
|
|
if subscribed_fields != "messages,messaging_postbacks":
|
|
return JSONResponse(
|
|
{"error": {"message": "invalid subscribed_fields", "code": 100}}, status_code=400,
|
|
)
|
|
if hook_hits is not None:
|
|
hook_hits.append({
|
|
"meta_page_subscription": page_id,
|
|
"subscribed_fields": subscribed_fields,
|
|
})
|
|
return {"success": True}
|
|
|
|
@up.get("/auth.test")
|
|
async def slack_auth_test(request: Request):
|
|
# Faithful Slack stand-in: it answers HTTP 200 even for a DEAD token and signals failure
|
|
# only via {"ok": false}. Checking the status alone would happily accept a bad token.
|
|
# It also reports the token's scopes in a response HEADER, not the body.
|
|
if "good" in request.headers.get("authorization", ""):
|
|
return JSONResponse(
|
|
{"ok": True, "team": "Acme Workspace", "team_id": "T0ACME", "user": "treg"},
|
|
headers={"x-oauth-scopes": "chat:write,channels:read,users:read"},
|
|
)
|
|
return JSONResponse({"ok": False, "error": "invalid_auth"})
|
|
|
|
@up.post("/hook")
|
|
async def hook(request: Request) -> dict:
|
|
# records health webhook POSTs so alerting tests can assert the webhook actually fired.
|
|
if hook_hits is not None:
|
|
hook_hits.append(await request.json())
|
|
return {"ok": True}
|
|
|
|
@up.get("/units")
|
|
async def semrush_units():
|
|
# Semrush's free unit-balance check answers HTTP 200 with a PLAIN-TEXT body, not JSON — the
|
|
# key-connect probe must not try to JSON-parse it, or a valid key reads as "unreachable".
|
|
from fastapi.responses import PlainTextResponse
|
|
return PlainTextResponse("API units balance: 1200")
|
|
|
|
@up.get("/units-bad")
|
|
async def semrush_units_bad():
|
|
# Semrush signals a bad key with HTTP 200 and a text body like "ERROR 120 :: ...", so the
|
|
# probe must read the body, not the status, to reject it.
|
|
from fastapi.responses import PlainTextResponse
|
|
return PlainTextResponse("ERROR 120 :: wrong key")
|
|
|
|
@up.get("/verify-field")
|
|
async def verify_field(request: Request):
|
|
# Emulates Apollo: HTTP 200 even for a BAD key, validity signalled only by a body field
|
|
# (is_logged_in). A key containing "good" is valid. The probe must read the field, not status.
|
|
ok = "good" in request.headers.get("x-api-key", "")
|
|
return {"healthy": True, "is_logged_in": ok}
|
|
|
|
@up.get("/credit-json-as-text")
|
|
async def credit_json_as_text(request: Request):
|
|
# Emulates ScrapeCreators: a real JSON body served with a text/plain content-type. The probe
|
|
# must parse it anyway to read token_verify_field, not gate on the mislabelled header.
|
|
from fastapi.responses import PlainTextResponse
|
|
return PlainTextResponse('{"success":true,"creditCount":17220}')
|
|
|
|
@up.get("/needs-query")
|
|
async def needs_query(request: Request):
|
|
# Emulates PDL/Akta/JustOneAPI/SpyFu: a probe_path with a required ?query. httpx drops a URL's
|
|
# own query when params= is passed, so a valid key 400'd until the query was merged into params.
|
|
if not request.query_params.get("field"):
|
|
from fastapi.responses import JSONResponse
|
|
return JSONResponse({"message": "field is required"}, status_code=400)
|
|
return {"ok": True}
|
|
|
|
@up.get("/requires-version")
|
|
async def requires_version(request: Request):
|
|
# Crustdata requires this protocol header on every route, including its free credential
|
|
# probe. The provisioner must stamp it without asking each caller to remember it.
|
|
if request.headers.get("x-api-version") != "2025-11-01":
|
|
from fastapi.responses import JSONResponse
|
|
return JSONResponse({"message": "x-api-version is required"}, status_code=400)
|
|
return {"ok": True, "version": request.headers["x-api-version"]}
|
|
|
|
@up.api_route("/{path:path}", methods=["GET", "POST", "PUT", "PATCH", "DELETE"])
|
|
async def echo(request: Request) -> dict:
|
|
body = (await request.body()).decode()
|
|
return {
|
|
"auth": request.headers.get("authorization"),
|
|
"headers": {k.lower(): v for k, v in request.headers.items()},
|
|
"query": dict(request.query_params),
|
|
"query_multi": request.query_params.multi_items(), # preserves duplicate keys
|
|
"body": body,
|
|
"raw_path": request.scope.get("raw_path", b"").decode(), # pre-decode bytes, for %2f fidelity asserts
|
|
}
|
|
|
|
return up
|
|
|
|
|
|
async def verified_identity(client, email):
|
|
"""Real OTP proof, capturing mail delivery even when Postgres hides dev codes."""
|
|
from unittest.mock import patch
|
|
import httpx
|
|
from treg import email as email_sender
|
|
|
|
delivered = {}
|
|
|
|
async def receive(email, code, **kwargs):
|
|
delivered[email] = code
|
|
|
|
previous_cookies = httpx.Cookies(client.cookies)
|
|
try:
|
|
with patch.object(email_sender, "send_otp", receive):
|
|
started = await client.post("/auth/email/start", json={"email": email})
|
|
assert started.status_code == 200, started.text
|
|
code = started.json().get("dev_code") or delivered[email]
|
|
proof = await client.post("/auth/email/verify", json={"email": email, "code": code})
|
|
assert proof.status_code == 200, proof.text
|
|
return proof.json()["token"]
|
|
finally:
|
|
client.cookies = previous_cookies
|
|
|
|
|
|
async def verified_signup(client, *, json, headers=None):
|
|
"""Funded test identity through OTP and team creation, with fixture identity fields."""
|
|
import httpx
|
|
from sqlmodel import select
|
|
from treg.infra.db import session_maker
|
|
from treg.models import User
|
|
|
|
email = json["email"]
|
|
token = await verified_identity(client, email)
|
|
response = await client.post("/orgs", json={"name": email},
|
|
headers={**(headers or {}), "X-Treg-Token": token})
|
|
if response.status_code != 200:
|
|
return response
|
|
async with session_maker() as db:
|
|
user = (await db.execute(select(User).where(User.email == email))).scalar_one()
|
|
user_id = user.id
|
|
return httpx.Response(response.status_code,
|
|
json={**response.json(), "id": user_id, "email": email}, request=response.request)
|
|
|
|
|
|
async def funded_user(client, email, *, micro=1_000_000):
|
|
"""A per-org token whose team holds money.
|
|
|
|
`POST /users` is legacy registration: the user it mints is UNVERIFIED, and the signup credit is
|
|
now verified-only (`claim_signup_promo`), so that team starts at zero. A test that has to PAY
|
|
for something funds the team here instead of leaning on a promo that no longer arrives.
|
|
Returns the whole registration body, so `["token"]` and `["org_id"]` both work.
|
|
"""
|
|
from treg.domain import money
|
|
from treg.infra.db import session_maker
|
|
|
|
body = (await client.post("/users", json={"email": email})).json()
|
|
async with session_maker() as db:
|
|
await money.grant(db, body["org_id"], amount_micro=micro, kind="promotional")
|
|
await db.commit()
|
|
return body
|
|
|
|
|
|
async def drain_background_writes():
|
|
# Postgres needs a session-scoped event loop so asyncpg can safely pool connections. That also
|
|
# lets fire-and-forget audit writes survive between tests, so drain both sides of reset_db():
|
|
# before it, to keep an old write out of the new schema, and after the test, to finish its own.
|
|
await audit.drain()
|
|
# Same discipline for the archive's fire-and-forget recordings: a still-open recording
|
|
# transaction from the PREVIOUS test blocks reset_db's DROP TABLE on Postgres (sqlite
|
|
# forgives it) — the serial CI job hung exactly here, 5-minute faulthandler timeouts on
|
|
# whichever archive test ran next (2026-08-28, twice).
|
|
await archive.drain()
|
|
# And the managed-key last-used writer: ids restart with every reset_db(), so a write or a
|
|
# throttle claim left over from one test would land on, or suppress, the next test's key.
|
|
await managed_keys.drain_last_used()
|
|
managed_keys._last_used_claims.clear()
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
async def _drain_around_each_test():
|
|
"""Every test starts and ends with no background write in flight. A fixture that calls
|
|
reset_db() itself (many do) would otherwise race the previous test's audit rows: a call record
|
|
committed between reset_db's delete of `callrecord` and its delete of `org` breaks the foreign
|
|
key, and one that lands after the reset attaches to the next test's rewound ids. Autouse
|
|
fixtures set up first and tear down last, so this brackets every other fixture's reset."""
|
|
await drain_background_writes()
|
|
yield
|
|
await drain_background_writes()
|
|
|
|
|
|
@pytest.fixture
|
|
async def clients():
|
|
# The archive report's 30s server-side cache would outlive this reset and serve the previous
|
|
# test's numbers — clear it with the schema.
|
|
from treg.routers import admin as admin_routes
|
|
admin_routes._archive_report_cache.clear()
|
|
await reset_db()
|
|
await app.state.endpoint_observation_reader.reset()
|
|
app.state.hook_hits = [] # webhook POSTs the upstream received (for alerting assertions)
|
|
app.state.http = AsyncClient(transport=ASGITransport(app=make_upstream(app.state.hook_hits)), base_url="http://upstream")
|
|
try:
|
|
async with AsyncClient(transport=ASGITransport(app=app), base_url="http://registry") as c:
|
|
r = await verified_signup(c, json={"email": "tim@superdesign.dev"})
|
|
assert r.status_code == 200, r.text
|
|
c.headers["X-Treg-Token"] = r.json()["token"] # authed by default from here on
|
|
yield c
|
|
finally:
|
|
await drain_background_writes()
|
|
await app.state.http.aclose()
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _reset_call_path_caches():
|
|
"""The call path keeps in-process copies of ratestore state (the capacity view: 'provider X is
|
|
exhausted' for 60 s), the overflow route view and per-provider rate buckets. `reset_db()` wipes
|
|
the tables, not process memory — so a test that relays a vendor 402 would otherwise leave the
|
|
NEXT test's tier-4 call refused with a 503 it never asked for (CI, xdist worker gw3, 2026-08-28).
|
|
Each cache is optional: the modules land in successive PRs of the capacity stack."""
|
|
def _clear() -> None:
|
|
try:
|
|
from treg.domain.capacity.view import view as capacity_view
|
|
capacity_view.invalidate()
|
|
capacity_view._states = {}
|
|
capacity_view._locks = {}
|
|
from treg.domain.capacity import marks as capacity_marks
|
|
capacity_marks._last_probe.clear()
|
|
except ImportError:
|
|
pass
|
|
try:
|
|
from treg.domain.capacity.routes_view import view as routes_view
|
|
routes_view.invalidate()
|
|
routes_view._routes = []
|
|
except ImportError:
|
|
pass
|
|
try:
|
|
from treg.infra.upstream.limiter import limiter
|
|
limiter.reset()
|
|
except ImportError:
|
|
pass
|
|
# The shared store's in-process fallback (the review-invitation budget): org ids restart
|
|
# with every reset_db(), so a counter left over would ration the NEXT test's team.
|
|
from treg.infra import kv
|
|
kv._store = None
|
|
# Who a key is, for the hub's lists: remembered a minute, and keys repeat across resets.
|
|
from treg.routers import hub_gate
|
|
hub_gate._readers.clear()
|
|
# The catalog's agent verdicts: a five-minute fold of callreview, which reset_db() empties.
|
|
from treg.application import feedback
|
|
feedback.forget_endpoint_verdicts()
|
|
_clear()
|
|
yield
|
|
_clear()
|
|
|
|
|
|
@pytest.fixture
|
|
def posthog_events(monkeypatch):
|
|
"""The PostHog mirror, switched on for one test with its POST stubbed. `await posthog_events()`
|
|
drains the queue and returns every `tool_called` event the server would have batched out, so a
|
|
test can assert on the product-analytics shape of a call without the two-second flush timer."""
|
|
from treg import analytics
|
|
from treg.config import get_settings
|
|
|
|
monkeypatch.setattr(get_settings(), "posthog_key", "phc_test_suite", raising=False)
|
|
sent: list[dict] = []
|
|
|
|
async def _keep(batch):
|
|
sent.extend(batch)
|
|
|
|
monkeypatch.setattr(analytics, "_post", _keep)
|
|
analytics._queue.clear()
|
|
|
|
async def collect(event: str = "tool_called") -> list[dict]:
|
|
await analytics.drain()
|
|
return [e for e in sent if e["event"] == event]
|
|
|
|
yield collect
|
|
analytics._queue.clear()
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _no_ambient_treg_identity(monkeypatch):
|
|
"""A dev machine may carry a per-agent identity in its environment (TREG_TOKEN et al — the
|
|
'Scope this agent' setup persists them into the coding agent's global env, and the CLI lets
|
|
them beat any config). The suite must not change behavior because of who is running it."""
|
|
for var in ("TREG_TOKEN", "TREG_ORG", "TREG_URL", "TREG_CLIENT"):
|
|
monkeypatch.delenv(var, raising=False)
|
|
|
|
|
|
@pytest.fixture
|
|
def kitt_on(monkeypatch):
|
|
from treg.config import get_settings
|
|
monkeypatch.setenv('TREG_PLATFORM_KEY_TRYKITT', 'TEST-KITT-KEY')
|
|
monkeypatch.setenv('TREG_PLATFORM_PROVIDERS', 'trykitt')
|
|
get_settings.cache_clear()
|
|
yield
|
|
get_settings.cache_clear()
|
|
|
|
|
|
# ---- ContactOut ----
|
|
|
|
@pytest.fixture
|
|
def contactout_platform(monkeypatch):
|
|
monkeypatch.setenv("TREG_PLATFORM_KEY_CONTACTOUT", "PLATFORM-TEST")
|
|
monkeypatch.setenv("TREG_PLATFORM_PROVIDERS", "contactout")
|
|
get_settings.cache_clear()
|
|
yield
|
|
get_settings.cache_clear()
|
|
|
|
|
|
_SESSION_CWD = os.getcwd()
|
|
|
|
|
|
@pytest.hookimpl(wrapper=True)
|
|
def pytest_runtest_teardown(item, nextitem):
|
|
"""Every test ends in the directory the run started in. A test that changes it (a bare
|
|
`os.chdir`, or product code that chdirs while a stub stands in for `exec`) makes later tests
|
|
fail on relative paths, far from the cause and only in some orders. Runs after every fixture
|
|
has torn down, so `monkeypatch.chdir` restores first; the leaking test errors here instead."""
|
|
result = yield
|
|
if os.getcwd() != _SESSION_CWD:
|
|
leaked = os.getcwd()
|
|
os.chdir(_SESSION_CWD)
|
|
raise AssertionError(f"{item.nodeid} left the working directory at {leaked}; "
|
|
"use monkeypatch.chdir so it is restored")
|
|
return result
|