diff --git a/docs/context/ops/capacity.md b/docs/context/ops/capacity.md index d6b0a55e..4a669242 100644 --- a/docs/context/ops/capacity.md +++ b/docs/context/ops/capacity.md @@ -397,6 +397,8 @@ pays the aggregator's real price, 0% markup, disclosed in-band when it ships (st - **`infra/upstream/aggregators/`** — the envelopes, and nothing else: `build()` wraps the vendor request (Orthogonal `POST /run {api, path, query, body}`; Monid `POST /run {provider, endpoint, input}`), `parse()` unwraps the vendor status + body + the real in-band charge, and + `build()` restores JSON scalar types that Monid validates and converts Akta enrichment's native + comma-separated `sections` query value to Monid's array-shaped envelope field, and names who to blame when the aggregator itself refused (`AGGREGATOR_SIDE` = `aggregator_auth`, `aggregator_balance`, `malformed` - the call path marks the aggregator unhealthy for everyone, the verifier leaves the route alone; `contract` - the aggregator's own per-request refusal, including @@ -421,7 +423,9 @@ pays the aggregator's real price, 0% markup, disclosed in-band when it ships (st - **`verify.py`** + `treg-worker overflow verify` — the weekly re-verify: one cheap call per route through the aggregator (and, when we hold the vendor key, directly), compare the shape fingerprint (keys and list/leaf markers, values ignored), stamp `last_verified_at` or disable - with the reason. Two per-route price caps and one run budget: a route that is enabled or was + with the reason. A mismatch prints a bounded key-only structural diff, never response values, so + an operator can distinguish omitted metadata from an incompatible body without exposing PII. + Two per-route price caps and one run budget: a route that is enabled or was stamped before is a **renewal**, held to `--renew-max-usd` (default $1); a never-verified pair is **discovery**, visited only under `--all` and held to `--max-usd` (default 2¢). Renewals go first, oldest stamp first, so the route nearest its 7-day decay is reached before `--budget-usd` diff --git a/src/treg/domain/capacity/overflow_seed.json b/src/treg/domain/capacity/overflow_seed.json index 11836bc1..74247222 100644 --- a/src/treg/domain/capacity/overflow_seed.json +++ b/src/treg/domain/capacity/overflow_seed.json @@ -11,7 +11,7 @@ "path": "/v1/company/enrichment/", "matched_at": "2026-08-26", "verified_at": null, - "single_result": false + "single_result": true }, { "aggregator": "monid", diff --git a/src/treg/domain/capacity/verify.py b/src/treg/domain/capacity/verify.py index 316910d3..7481d7f2 100644 --- a/src/treg/domain/capacity/verify.py +++ b/src/treg/domain/capacity/verify.py @@ -38,6 +38,30 @@ def shapes_match(a: bytes, b: bytes) -> bool | None: return None +def _shape_paths(value, path: str = "$") -> set[str]: + """PII-free structural paths for explaining a failed shape comparison.""" + if isinstance(value, dict): + paths = {f"{path}:object"} + for key, child in value.items(): + paths |= _shape_paths(child, f"{path}.{key}") + return paths + if isinstance(value, list): + return {f"{path}:list"} | (_shape_paths(value[0], f"{path}[]") if value else set()) + return {f"{path}:leaf"} + + +def shape_difference(a: bytes, b: bytes) -> str: + """Describe only structural differences; never include response values.""" + try: + direct = _shape_paths(json.loads(a)) + relay = _shape_paths(json.loads(b)) + except ValueError: + return "non-JSON response" + direct_only = sorted(direct - relay)[:8] + relay_only = sorted(relay - direct)[:8] + return f"direct-only={direct_only or '-'}; relay-only={relay_only or '-'}"[:500] + + @dataclass class Verification: endpoint_id: str @@ -145,4 +169,5 @@ async def verify_route(client: httpx.AsyncClient, route, *, key: str, direct: tu same = shapes_match(dr.content, res.upstream_body) return Verification(route.endpoint_id, route.aggregator, dr.status_code, res.upstream_status, same, res.cost_micro, now if same else None, - note="" if same else f"direct {dr.status_code}, relay {res.upstream_status}, shape differs") + note=("" if same else f"direct {dr.status_code}, relay {res.upstream_status}, " + f"shape differs: {shape_difference(dr.content, res.upstream_body)}")) diff --git a/src/treg/infra/upstream/aggregators/monid.py b/src/treg/infra/upstream/aggregators/monid.py index b618eec2..063e9c25 100644 --- a/src/treg/infra/upstream/aggregators/monid.py +++ b/src/treg/infra/upstream/aggregators/monid.py @@ -22,6 +22,13 @@ NAME = "monid" _INT = re.compile(r"-?\d{1,15}") _FLOAT = re.compile(r"-?\d+\.\d+") +# Monid validates these query parameters against an array schema even though the vendor's native +# GET API accepts one comma-separated query value. This is envelope adaptation only: the same +# values reach the same vendor parameter, with no provider response modeling. +_ARRAY_QUERY_PARAMS = { + ("akta", "/v1/company/enrichment"): frozenset({"sections"}), +} + def _typed(v): """Monid validates `input` against the vendor's JSON schema, so a numeric query value must be a @@ -39,9 +46,22 @@ def _typed(v): return v +def _query_params(route, query: list[tuple[str, str]] | dict) -> dict: + pairs = list(query.items() if isinstance(query, dict) else query) + arrays = _ARRAY_QUERY_PARAMS.get((route.agg_slug, route.agg_path.rstrip("/")), frozenset()) + out: dict = {} + for key, value in pairs: + if key not in arrays: + out[key] = _typed(value) + continue + values = value if isinstance(value, list) else str(value).split(",") + out.setdefault(key, []).extend(_typed(v.strip()) for v in values if str(v).strip()) + return out + + def build(route, key: str, query: list[tuple[str, str]] | dict, body: bytes | None, path_params: dict | None = None, *, params_as_body: bool = False) -> AggregatorRequest: - items = {k: _typed(v) for k, v in (query.items() if isinstance(query, dict) else query)} + items = _query_params(route, query) parsed: dict = {} if body: try: diff --git a/tests/test_capacity_overflow_routes.py b/tests/test_capacity_overflow_routes.py index ffc4bac0..5ccaa0fa 100644 --- a/tests/test_capacity_overflow_routes.py +++ b/tests/test_capacity_overflow_routes.py @@ -159,6 +159,21 @@ async def test_verified_akta_news_monid_fallback_is_enabled_at_the_observed_defa assert row.ratio == 1 +async def test_akta_enrichment_monid_route_is_eligible_only_after_live_verification(): + await reset_db() + candidate = next(x for x in R.load_seed() + if x["endpoint_id"] == "akta.companies.enrich" and x["aggregator"] == "monid") + assert candidate["verified_at"] is None and candidate["single_result"] is True + candidate = {**candidate, "verified_at": "2026-09-25"} + async with session_maker() as db: + await ensure_policies(db, has_key=lambda p: True) + await R.apply_sync(db, [candidate], catalog=catalog_store.load(), + now=R._dt("2026-09-25T12:00:00")) + await db.commit() + row = await db.get(OverflowRoute, ("akta.companies.enrich", "monid")) + assert row.enabled and row.agg_unit == "result" and row.single_result is True + + def test_route_for_orders_orthogonal_first(): rs = [_route(aggregator="monid", enabled=True), _route(aggregator="orthogonal", enabled=True), _route(aggregator="orthogonal", endpoint_id="x", enabled=True), _route(aggregator="monid", enabled=False, endpoint_id="y")] @@ -363,6 +378,11 @@ def test_monid_build_and_parse_fixtures(): assert req.json == {"provider": "hunterio", "endpoint": "/domain-search", "input": {"queryParams": {"domain": "stripe.com", "limit": 1, "score": 0.5, "raw": True, "id": "007a"}, "body": {}, "pathParams": {}}}, "Monid validates JSON types: numeric strings become numbers (live 2026-08-28)" + akta = _route(aggregator="monid", agg_slug="akta", agg_path="/v1/company/enrichment") + req = monid.build(akta, "K", {"company": "canva.com", "sections": "location,technology"}, None) + assert req.json["input"]["queryParams"] == { + "company": "canva.com", "sections": ["location", "technology"], + }, "Monid's Akta schema models the vendor's comma-separated sections parameter as an array" alt = monid.build(r, "K", {"domain": "stripe.com"}, None, params_as_body=True) assert alt.json["input"] == {"queryParams": {}, "body": {"domain": "stripe.com"}, "pathParams": {}} ok = _fixture("monid_ok_sync") @@ -445,6 +465,9 @@ def test_shape_fingerprint_ignores_values_but_not_structure(): b = b'{"data":{"email":"b@y.io","score":1,"sources":[{"uri":"v"},{"uri":"w"}]}}' c = b'{"data":{"email":"b@y.io"}}' assert V.shapes_match(a, b) is True and V.shapes_match(a, c) is False and V.shapes_match(a, b"nope") is None + diff = V.shape_difference(a, c) + assert "$.data.score:leaf" in diff and "$.data.sources:list" in diff + assert "a@x.io" not in diff def test_shape_empty_vs_nonempty_list_differs_but_both_empty_match():