mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
fix: stop billing empty answers, make /admin/errors read-only, add treg --json call (#675)
* fix(billing): stop charging for empty and company-blind answers Per-result searches settled at the requested page size even when the vendor returned nothing: count the billed rows for CompanyEnrich (2-credit minimum), Findymail employees, Icypeas bulk (FOUND only), TheCompaniesAPI search (and simplified=true, which is free) and Serpstat (error envelopes free, 1-credit minimum on an empty result). Icypeas profile-URL misses settle at zero through an expect rule. Routing: Prospeo 400 NO_RESULTS is a declared miss, not a caller fault; the CompanyEnrich scroll adapter is removed so one empty question is not routed twice; people.search declares scoping: [company_domain], so a title-only provider is dropped from a company-scoped plan instead of billing title-matched strangers once the others miss. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMPFwsMRowa1X7kXc5nZyw * fix(admin): make GET /admin/errors read-only; purge evidence from a worker Loading the errors page ran a platform-wide UPDATE blanking failed-call evidence past the 14-day window. The purge moves to application/evidence_retention.py behind 'treg-worker admin purge-evidence', in bounded id batches. The view withholds evidence past the window itself, so an unscheduled purge never widens what it shows. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMPFwsMRowa1X7kXc5nZyw * feat(cli): treg --json call prints one parseable envelope With --json, call prints {"result": <body>, "_treg": {http_status, call_id, charged_micro|reserved_micro, ...}} on one line and nothing on stderr, so a script that merges the streams still parses every answer. Default output and --await are unchanged. Agent docs say to use it when parsing, and to check a few results before looping. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMPFwsMRowa1X7kXc5nZyw * fix(billing): price counted rows at their credits and cap them at the hold For a credit-priced per_result row the frozen unit is one provider credit, so counting rows alone under-billed CompanyEnrich (2 credits a person, and its 2-credit minimum) and Icypeas reverse-email bulk (10 a hit). Rows whose catalog unit names an input entity reserve per thing asked about, so a row count is capped at the hold: counting may lower a bill, never raise it. simplified=true is free only on TheCompaniesAPI endpoints that declare it. Also: a malformed contract scoping list fails the catalog load; --json call keeps stderr silent on the WAF base64 retry; AGENTS.md records the evidence-retention purge as the second callrecord writer. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMPFwsMRowa1X7kXc5nZyw * fix(billing): bill Icypeas company scrapes at the company rate; never exceed a zero hold Also corrects the AGENTS.md callrecord-writer note and the scoping comment (a bad list empties routing rather than failing the catalog). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TMPFwsMRowa1X7kXc5nZyw --------- Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
9d9d7f6996
commit
fd5a2c0033
@@ -203,6 +203,7 @@ Regenerate via `scripts/build-map.py`.
|
||||
| `src/treg/application/catalog_find.py` | architecture/search-experiment.md |
|
||||
| `src/treg/application/catalog_stats.py` | architecture/catalog.md |
|
||||
| `src/treg/application/connect.py` | architecture/auth-secrets.md, architecture/composition.md, guides/expanding-a-category.md, interface/api.md |
|
||||
| `src/treg/application/evidence_retention.py` | architecture/super-admin.md |
|
||||
| `src/treg/application/feedback.py` | architecture/feedback.md |
|
||||
| `src/treg/application/jev_xboost.py` | interface/seo.md |
|
||||
| `src/treg/application/media.py` | architecture/media.md |
|
||||
@@ -592,7 +593,7 @@ Regenerate via `scripts/build-map.py`.
|
||||
| `architecture/multi-tenancy.md` | `access.py`, `0042_pinned_read_scope.py`, `test_pinned_read_scope.py`, `models.py`, `api.py`, `caller_metadata.py`, `auth.py`, `asynctasks.py`, `resolve.py`, `provider_resources.py`, `provider_resources.py`, `provider_resources.py`, `0043_provider_resources.py`, `signup.py`, `access.py`, `budgets.py`, `publicdemo.py`, `teams.py`, `usage.py`, `access.py`, `api_keys.py`, `session.py`, `promotions.py`, `test_team_limit.py`, `test_auth.py`, `test_token_revocation.py`, `auth.py`, `orgs.py`, `resources.py`, `bundles.py`, `db.py`, `0017_async_task_record.py`, `0018_async_resource_ownership.py`, `test_router_dependencies.py`, `test_asynctasks.py` |
|
||||
| `architecture/proxy-model.md` | `relay.py`, `ssrf.py`, `api.py`, `authorize.py`, `idempotency.py`, `intake.py`, `resolve.py`, `reserve.py`, `settle.py`, `evidence.py`, `service.py`, `types.py`, `asynctasks.py`, `client_identity.py`, `call_surface.py`, `sandbox_identity.py`, `access.py`, `publicdemo.py`, `usage.py`, `call.py`, `test_ssrf_public_addresses.py`, `test_call_application_contract.py`, `test_call_cancellation.py`, `test_call_response_limits.py`, `test_error_capture.py`, `test_marketplace_call.py`, `test_oauth_billed.py`, `test_passthrough.py`, `test_tag_billing.py`, `test_tag_billing_adversarial.py`, `test_call_architecture.py`, `test_asynctasks.py`, `test_relay_content_length.py` |
|
||||
| `architecture/search-experiment.md` | `search_experiment.py`, `interleave.py`, `judge.py`, `0041_searchlog.py`, `search_experiment_report.sql`, `test_search_experiment.py`, `catalog_find.py`, `test_catalog_find.py` |
|
||||
| `architecture/super-admin.md` | `api.py`, `admin.py`, `access.py`, `config.py` |
|
||||
| `architecture/super-admin.md` | `api.py`, `admin.py`, `evidence_retention.py`, `access.py`, `config.py` |
|
||||
| `foundation/charter.md` | `2026-06-30-jason-tools-registry.md`, `README.md` |
|
||||
| `guides/expanding-a-category.md` | `oauth_providers.py`, `authorization.py`, `oauth_flow.py`, `oauth_exchange.py`, `connect.py`, `connections.py`, `config.py` |
|
||||
| `interface/api.md` | `media.py`, `sitetrack.js`, `api.py`, `bootstrap_handlers.py`, `bootstrap_http.py`, `call_surface.py`, `caller_metadata.py`, `client_identity.py`, `auth.py`, `provider_resources.py`, `access.py`, `authorize.py`, `idempotency.py`, `intake.py`, `resolve.py`, `reserve.py`, `settle.py`, `evidence.py`, `service.py`, `types.py`, `relay.py`, `connect.py`, `onboard.py`, `referrals.py`, `signup.py`, `__init__.py`, `admin.py`, `auth.py`, `auth_helpers.py`, `billing.py`, `call.py`, `catalog.py`, `connections.py`, `onboard.py`, `orgs.py`, `provider_resources.py`, `api_keys.py`, `resources.py`, `referrals.py`, `signup_cookies.py`, `web.py`, `access.py`, `api_keys.py`, `teams.py`, `access.py`, `budgets.py`, `publicdemo.py`, `usage.py`, `mcp_oauth.py`, `session.py`, `timeutil.py`, `store.py`, `email.py`, `runner.py`, `ratestore.py` |
|
||||
|
||||
@@ -105,7 +105,8 @@ agents then built against a constitution that was wrong.
|
||||
- **Table ownership.** One writer module per table; cross-domain reads are fine. Three recorded
|
||||
exceptions: only money writes `org.balance_micro`, the daily-spend counter (`spent_today_*`) and
|
||||
the auto-top-up fields; the call runtime may persist an OAuth token refresh into `secret`; audit
|
||||
writes `callrecord`, domains only read it.
|
||||
writes `callrecord`, domains only read it (`application/evidence_retention.py` also updates it,
|
||||
blanking the two evidence columns past retention).
|
||||
- **Feedback handling.** This repo owns `FeedbackHandling` and `FeedbackHandlingEvent` models and
|
||||
migrations; the private admin service is their only runtime writer. Original reports remain
|
||||
owned by the feedback domain. See `docs/context/architecture/feedback.md`.
|
||||
|
||||
@@ -137,6 +137,7 @@ treg tool add google-ads --base-url https://googleads.googleapis.com \
|
||||
| `treg catalog search` | `"what you want to do"` | find endpoints by capability |
|
||||
| `treg catalog get` | `ENDPOINT_ID` | docs, parameters, **the price**, and how you would be served |
|
||||
| `treg call ENDPOINT_ID` | `--query K=V`, `--data STR` | call it |
|
||||
| `treg --json call ENDPOINT_ID` | same | for scripts: one JSON line `{"result": <body>, "_treg": {http_status, call_id, charged_micro}}` on stdout, nothing on stderr |
|
||||
| `treg call ENDPOINT_ID --await` | `--timeout N` (default 900) | a generation call (video/image): submit, poll the provider, print the final response |
|
||||
| `treg host FILE` | `--content-type T`, `--json` | host a reference image/audio/video at a public URL a vendor can fetch (30 MB, 7-day TTL, free); prints the URL for `image_urls` / `audio_urls` |
|
||||
| `treg catalog request` | `"what's missing"` | searched, not there? file it — requests steer what gets added next |
|
||||
|
||||
@@ -33,7 +33,7 @@ covers (frontmatter `sources:`). Regenerate this index with
|
||||
| [Multi-tenancy — orgs, memberships, invites, per-org scoping](architecture/multi-tenancy.md) | shipped | access.py, 0042_pinned_read_scope.py, test_pinned_read_scope.py, models.py, … |
|
||||
| [The proxy — faithful credential-injecting relay + tool resolution](architecture/proxy-model.md) | shipped | relay.py, ssrf.py, api.py, authorize.py, … |
|
||||
| [Discovery experiment — a relevance judge behind catalog search, measured on what the caller does next](architecture/search-experiment.md) | building | search_experiment.py, interleave.py, judge.py, 0041_searchlog.py, … |
|
||||
| [Super-admin — cross-tenant read + control](architecture/super-admin.md) | shipped | api.py, admin.py, access.py, config.py |
|
||||
| [Super-admin — cross-tenant read + control](architecture/super-admin.md) | shipped | api.py, admin.py, evidence_retention.py, access.py, … |
|
||||
|
||||
## Interfaces (API · CLI · skill)
|
||||
|
||||
|
||||
@@ -1708,6 +1708,11 @@ to choose (`docs/CAPABILITY-ROUTING-PLAN.md`). Everything else in the catalog st
|
||||
admission-only contract: its adapters verify like any other (which is what the archive's
|
||||
`has_result_rules` reads), but no `treg.<capability>` row is ever generated from it. For a
|
||||
capability whose "children" are one provider's price tiers, not a choice treg should make.
|
||||
`scoping` names identity keys that scope the answer rather than describe it (`people.search`:
|
||||
`company_domain`). A candidate whose adapter never sends one the caller supplied is dropped from
|
||||
the plan with the reason, not ranked down like an ignored filter: a title-only search asked for
|
||||
one company's CEO returns title-matched strangers for any company and bills them as a hit. The
|
||||
rule is per candidate, so `{q, company_domain}` also drops the `q`-only providers.
|
||||
- **Adapters** — `adapters.yaml`, one per endpoint: `accepts` (identity variants), `in` (contract
|
||||
field → `queryParams.x` / `body.x`), `const` (fixed provider params), `out` (core field →
|
||||
expression over the body), `miss`. The expression language (`domain/catalog/routing/paths.py`)
|
||||
|
||||
@@ -282,8 +282,9 @@ uses this metadata, never the encrypted token's shape.
|
||||
`error_request` / `error_response` hold redacted, truncated failure evidence across platform,
|
||||
own-key and own-tool calls. Successes leave them empty. Captured provider headers use an
|
||||
allowlist covering retry/auth/rate-limit and request/trace identifiers. `/calls` neither fetches
|
||||
nor exposes these wide fields; `GET /admin/errors` owns access and the 14-day retention purge
|
||||
(replacing expired evidence with `<expired>`).
|
||||
nor exposes these wide fields; `GET /admin/errors` owns read access (read-only) and the
|
||||
`treg-worker admin purge-evidence` cron (`application/evidence_retention.py`) the 14-day
|
||||
retention purge, replacing expired evidence with `<expired>`.
|
||||
|
||||
Redaction in `application.call.evidence` is security-sensitive:
|
||||
|
||||
@@ -297,7 +298,8 @@ uses this metadata, never the encrypted token's shape.
|
||||
|
||||
Unmetered uploads are buffered for evidence only with declared `Content-Length <= 64 KiB`.
|
||||
Failed streaming responses retain at most the first 8 KiB, replaying all bytes to the caller.
|
||||
Purging stays on the admin path; request-session dependencies do not commit a lazy purge marker.
|
||||
Purging never runs on a request: an admin page reading errors once blanked evidence
|
||||
platform-wide on every load.
|
||||
|
||||
`archive_key_hash` / `archive_content_hash` link eligible metered platform responses to
|
||||
`ArchiveKey` / `ArchiveSnapshot` for `GET /calls/{id}/result`. They are nullable, unindexed
|
||||
|
||||
@@ -4,6 +4,7 @@ status: shipped
|
||||
sources:
|
||||
- src/treg/api.py
|
||||
- src/treg/routers/admin.py
|
||||
- src/treg/application/evidence_retention.py
|
||||
- src/treg/domain/identity/access.py
|
||||
- src/treg/config.py
|
||||
related:
|
||||
@@ -60,10 +61,10 @@ endpoints are unaffected (they use `require_superadmin`).
|
||||
[data-model](data-model.md)). `tier` filters an exact marketplace tier; an empty value selects plain
|
||||
own tools. Superadmin and not org-admin
|
||||
because the rows hold customers' request content; `GET /calls` deliberately does **not** expose
|
||||
these columns, and it defers them so they are not even fetched. This route also performs the
|
||||
14-day retention pass (`_purge_expired_error_evidence`, blanking to `'<expired>'` on its own
|
||||
committed session) — ageing lives here because there is no scheduler and the request path cannot
|
||||
hold a lazy marker, `get_admin_session` never committing one.
|
||||
these columns, and it defers them so they are not even fetched. The route is read-only: the
|
||||
14-day retention purge (blanking both columns to `'<expired>'`) is the `treg-worker admin
|
||||
purge-evidence` cron (`application/evidence_retention.py`), in bounded batches. A row past the
|
||||
window is listed as `expired` with no evidence even before the cron reaches it.
|
||||
- **Reconciliation (Phase 5):** `admin_reconcile_drift|spend|repeats` (`?since_days=30`) — cross-org
|
||||
aggregates over platform-tier spend, so super-admin and not org-admin: price drift per endpoint,
|
||||
settled spend per provider (the invoice comparison), and the repeat-query rate. Query-time reports
|
||||
|
||||
@@ -327,8 +327,8 @@ validated before resolving the shared HTTP client. `/auth/logout` remains an HTT
|
||||
itself ran that journal count - 2.8 s per call for a member with 110k rows that day.
|
||||
- **Super-admin (cross-tenant, `require_superadmin`):** `/admin/stats|orgs|orgs/{id}|users|tools|calls|
|
||||
errors|health` (reads - `errors` is failed calls across every credential tier with captured,
|
||||
admin-only request/response evidence, supports a `tier` filter, and runs the 14-day retention pass;
|
||||
see [super-admin](../architecture/super-admin.md))
|
||||
admin-only request/response evidence, supports a `tier` filter, and withholds evidence past the
|
||||
14-day retention window, which the `treg-worker admin purge-evidence` cron blanks; see [super-admin](../architecture/super-admin.md))
|
||||
+ `/admin/users/{id}/superadmin|suspend`, `DELETE /admin/users/{id}`,
|
||||
`/admin/orgs/{id}/suspend`, `DELETE /admin/orgs/{id}` (Phase-2). See
|
||||
[super-admin](../architecture/super-admin.md).
|
||||
|
||||
@@ -75,7 +75,11 @@ epilog (a `mk()` helper + `_ex()` + `RawDescriptionHelpFormatter`), so `treg <cm
|
||||
`treg --version` / `treg version` print `cli_version()` (package metadata); `treg update` (`cmd_update`)
|
||||
re-runs the server's `install.sh` to upgrade the CLI in place. A global **`--json`** flag (stripped in
|
||||
`main` like `--org`) makes the human-table commands (`org ls`, `agents ls`, `catalog` in all its forms)
|
||||
emit raw JSON instead — one stable contract for agents; commands that already print JSON are unaffected.
|
||||
emit raw JSON instead — one stable contract for agents. On `call` (not `--await`) it prints one
|
||||
compact envelope, `{"result": <body>, "_treg": {http_status, call_id, charged_micro | reserved_micro,
|
||||
replay?, async?, hint?}}` (`_call_envelope`; text as a string, binary as base64), and suppresses the
|
||||
charge, hint and failure-diagnostic stderr lines: a script that merged the streams once discarded
|
||||
every result it had paid for. Exit status is unchanged (1 on HTTP >= 400).
|
||||
**`TREG_CONFIG`** points the CLI at an alternate config file (CI/agents/tests; default
|
||||
`~/.treg/config.json`). `org use` validates the slug against `/orgs`, then gets that membership's
|
||||
active Default key before it saves either value. If that exchange fails, the previous team and token
|
||||
|
||||
@@ -341,13 +341,17 @@ without importing the heavy database stack into the light `treg` CLI.
|
||||
- `treg-worker arena insights` folds new audit rows into the rolling Arena aggregate
|
||||
(`--max-seconds`, default 110, bounds one pass; schedule it every two minutes).
|
||||
- `treg-worker catalog stats` folds new audit rows into per-endpoint, per-day reliability buckets
|
||||
(`--max-rows`, default 500,000, bounds one pass; schedule it every few minutes). The catalog keeps
|
||||
computing observations live until this command has caught up once, so it can be scheduled after
|
||||
the application deploys, and a self-hosted registry that never schedules it loses nothing.
|
||||
- `treg-worker jev xboost` runs the `/jev` launch-radar demo once a day: it calls treg's own `/call/` API
|
||||
with `TREG_JEV_TREG_TOKEN` (a member token of the demo team, so the spend is an ordinary bill) and jev
|
||||
through the Vercel AI Gateway (`TREG_AI_GATEWAY_API_KEY`), and stores the run under Ephemeral for the page.
|
||||
Both variables also belong on the web service, which needs them for the visitor judge endpoint.
|
||||
(`--max-rows`, default 500,000, bounds one pass; schedule it every few minutes). The catalog keeps
|
||||
computing observations live until this command has caught up once, so it can be scheduled after
|
||||
the application deploys, and a self-hosted registry that never schedules it loses nothing.
|
||||
- `treg-worker admin purge-evidence` blanks failed-call evidence past the 14-day retention window
|
||||
(`--batch-size`, default 5000, rows per transaction; schedule it daily). `GET /admin/errors` is
|
||||
read-only and already withholds evidence past the window, so an unscheduled purge keeps the old
|
||||
bytes in the database but never shows them.
|
||||
|
||||
The two analytics commands exist so that no web process aggregates the audit table beside the
|
||||
money path; `callrecord` is read only through the persisted cursors they own. Workers call
|
||||
|
||||
@@ -127,6 +127,11 @@ Notes:
|
||||
to the balance — they take priority automatically). A 402 with `error: route_max_cost` is
|
||||
different: YOUR `X-Treg-Route-Max-Cost` header refused the call before anything was charged —
|
||||
ask for fewer rows/targets or raise the ceiling.
|
||||
- **Scripting many calls:** use `treg --json call …`. Stdout is one line,
|
||||
`{"result": <provider body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and nothing
|
||||
goes to stderr, so a script that merges the streams still parses every answer (`--await` output
|
||||
is unchanged). Run a handful and check the parsed results before looping over the whole list: a
|
||||
parse bug throws away answers that were already billed.
|
||||
- The real charge is the response header `X-Treg-Cost-Micro` (micro-USD), with `X-Treg-Call-Id`
|
||||
as the id to quote. On an asynchronous submission that header is the reserved ceiling; the CLI
|
||||
labels it as a reservation, and the terminal task settles the real charge. The catalog `~$/call`
|
||||
|
||||
@@ -130,6 +130,11 @@ Notes:
|
||||
to the balance — they take priority automatically). A 402 with `error: route_max_cost` is
|
||||
different: YOUR `X-Treg-Route-Max-Cost` header refused the call before anything was charged —
|
||||
ask for fewer rows/targets or raise the ceiling.
|
||||
- **Scripting many calls:** use `treg --json call …`. Stdout is one line,
|
||||
`{"result": <provider body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and nothing
|
||||
goes to stderr, so a script that merges the streams still parses every answer (`--await` output
|
||||
is unchanged). Run a handful and check the parsed results before looping over the whole list: a
|
||||
parse bug throws away answers that were already billed.
|
||||
- The real charge is the response header `X-Treg-Cost-Micro` (micro-USD), with `X-Treg-Call-Id`
|
||||
as the id to quote. On an asynchronous submission that header is the reserved ceiling; the CLI
|
||||
labels it as a reservation, and the terminal task settles the real charge. The catalog `~$/call`
|
||||
|
||||
@@ -111,6 +111,11 @@ Notes:
|
||||
to the balance — they take priority automatically). A 402 with `error: route_max_cost` is
|
||||
different: YOUR `X-Treg-Route-Max-Cost` header refused the call before anything was charged —
|
||||
ask for fewer rows/targets or raise the ceiling.
|
||||
- **Scripting many calls:** use `treg --json call …`. Stdout is one line,
|
||||
`{"result": <provider body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and nothing
|
||||
goes to stderr, so a script that merges the streams still parses every answer (`--await` output
|
||||
is unchanged). Run a handful and check the parsed results before looping over the whole list: a
|
||||
parse bug throws away answers that were already billed.
|
||||
- The real charge is the response header `X-Treg-Cost-Micro` (micro-USD), with `X-Treg-Call-Id`
|
||||
as the id to quote. On an asynchronous submission that header is the reserved ceiling; the CLI
|
||||
labels it as a reservation, and the terminal task settles the real charge. The catalog `~$/call`
|
||||
|
||||
@@ -125,6 +125,11 @@ Notes:
|
||||
to the balance — they take priority automatically). A 402 with `error: route_max_cost` is
|
||||
different: YOUR `X-Treg-Route-Max-Cost` header refused the call before anything was charged —
|
||||
ask for fewer rows/targets or raise the ceiling.
|
||||
- **Scripting many calls:** use `treg --json call …`. Stdout is one line,
|
||||
`{"result": <provider body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and nothing
|
||||
goes to stderr, so a script that merges the streams still parses every answer (`--await` output
|
||||
is unchanged). Run a handful and check the parsed results before looping over the whole list: a
|
||||
parse bug throws away answers that were already billed.
|
||||
- The real charge is the response header `X-Treg-Cost-Micro` (micro-USD), with `X-Treg-Call-Id`
|
||||
as the id to quote. On an asynchronous submission that header is the reserved ceiling; the CLI
|
||||
labels it as a reservation, and the terminal task settles the real charge. The catalog `~$/call`
|
||||
|
||||
@@ -109,6 +109,11 @@ Notes:
|
||||
to the balance — they take priority automatically). A 402 with `error: route_max_cost` is
|
||||
different: YOUR `X-Treg-Route-Max-Cost` header refused the call before anything was charged —
|
||||
ask for fewer rows/targets or raise the ceiling.
|
||||
- **Scripting many calls:** use `treg --json call …`. Stdout is one line,
|
||||
`{"result": <provider body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and nothing
|
||||
goes to stderr, so a script that merges the streams still parses every answer (`--await` output
|
||||
is unchanged). Run a handful and check the parsed results before looping over the whole list: a
|
||||
parse bug throws away answers that were already billed.
|
||||
- The real charge is the response header `X-Treg-Cost-Micro` (micro-USD), with `X-Treg-Call-Id`
|
||||
as the id to quote. On an asynchronous submission that header is the reserved ceiling; the CLI
|
||||
labels it as a reservation, and the terminal task settles the real charge. The catalog `~$/call`
|
||||
|
||||
@@ -36,7 +36,7 @@ from ...domain.catalog import stats as endpoint_stats
|
||||
from ...domain.catalog import store as catalog_store
|
||||
from ...domain.catalog.routing.contracts import canonical_identity, declared_miss, miss_status
|
||||
from ...domain.catalog.routing.plan import (
|
||||
MAX_ERROR_FALLBACKS, Candidate, Plan, candidates_for, cost_at, ignored_filters, rank,
|
||||
MAX_ERROR_FALLBACKS, Candidate, Plan, candidates_for, cost_at, ignored_filters, rank, unscoped,
|
||||
)
|
||||
from .. import asynctasks as async_task_app
|
||||
from . import async_bridge
|
||||
@@ -251,6 +251,14 @@ async def build_plan(ep: dict, identity_given: dict, caller, options: RouteOptio
|
||||
+ " | ".join("{" + ", ".join(v) + "}" for v in contract.identity),
|
||||
"variants": [list(v) for v in contract.identity]})
|
||||
raw, dropped = candidates_for(contract, cat.for_capability(ep["capability"]), cat.adapters, identity)
|
||||
scoped = []
|
||||
for e, ad, v in raw:
|
||||
if missing := unscoped(ad, contract, identity):
|
||||
dropped.append({"endpoint_id": e["id"], "why": f"cannot scope by {', '.join(missing)}; "
|
||||
"it would answer the same for any value"})
|
||||
else:
|
||||
scoped.append((e, ad, v))
|
||||
raw = scoped
|
||||
ids = [e["id"] for e, _, _ in raw]
|
||||
stats = await _observed_stats(ids)
|
||||
own: set[str] = set()
|
||||
|
||||
@@ -156,6 +156,80 @@ def _tavily_result_count(endpoint_id: str, doc: object) -> int | None:
|
||||
return None
|
||||
|
||||
|
||||
def _companyenrich_record_count(endpoint_id: str, doc: object) -> int | None:
|
||||
"""People returned by CompanyEnrich's search, floored at one: a person is 2 credits, and an
|
||||
empty page still costs the 2-credit minimum (catalog note, verified live). The reserve is the
|
||||
requested `pageSize`, so without counting an empty `{"items": []}` settled a whole page."""
|
||||
if endpoint_id not in (
|
||||
"companyenrich.people.search",
|
||||
"companyenrich.people.search.scroll",
|
||||
):
|
||||
return None
|
||||
if not isinstance(doc, dict):
|
||||
return None
|
||||
items = doc.get("items")
|
||||
if not isinstance(items, list):
|
||||
return None
|
||||
# 2 credits per person, minimum 1 unit charged (the 2-credit minimum on empty)
|
||||
return max(len(items), 1)
|
||||
|
||||
|
||||
def _icypeas_bulk_found_count(endpoint_id: str, doc: object) -> int | None:
|
||||
"""FOUND rows in an Icypeas bulk answer. Icypeas bills per found item and a NOT_FOUND row is
|
||||
free, while the reserve is the request's row count."""
|
||||
if endpoint_id not in (
|
||||
"icypeas.profile.url.bulk",
|
||||
"icypeas.people.identity.resolve.bulk",
|
||||
"icypeas.scrape.bulk",
|
||||
):
|
||||
return None
|
||||
if not isinstance(doc, dict):
|
||||
return None
|
||||
data = doc.get("data")
|
||||
if not isinstance(data, list):
|
||||
return None
|
||||
return sum(1 for item in data if isinstance(item, dict) and item.get("status") == "FOUND")
|
||||
|
||||
|
||||
def _serpstat_result_count(doc: object) -> int | None:
|
||||
"""Credits a Serpstat JSON-RPC answer bills, in rows. HTTP 200 carries both outcomes: an `error`
|
||||
envelope (bad token, exhausted limit, "Data not found") bills nothing; a `result` bills per row
|
||||
with the documented 1-credit minimum on an empty list. Rows live in `result.data[]`, or one
|
||||
level deeper for getKeywordTop (`result.data.top[]`). Any other shape (results keyed by the
|
||||
thing asked about) settles at the estimate rather than guessing a row count."""
|
||||
if not isinstance(doc, dict):
|
||||
return None
|
||||
if doc.get("error"):
|
||||
return 0
|
||||
result = doc.get("result")
|
||||
data = result.get("data") if isinstance(result, dict) else None
|
||||
if isinstance(data, dict):
|
||||
data = data.get("top")
|
||||
if isinstance(data, list):
|
||||
return max(len(data), 1)
|
||||
return None
|
||||
|
||||
|
||||
def _rows_billed_micro(mk: MarketplaceCall, ep: dict | None, rows: int | None,
|
||||
credits_per_row: Decimal | None = None) -> int | None:
|
||||
"""What `rows` billed rows cost, never more than the hold. For a credit-priced row
|
||||
`mk.unit_micro` is ONE provider credit, so it is scaled by the row's credits (`cost.value`:
|
||||
2 per CompanyEnrich person, 10 per Icypeas reverse-email hit). Capped at the reserve because a
|
||||
row whose catalog `unit` names an input entity (`call`, `keyword`, `domain`) reserves per thing
|
||||
asked about, not per row returned: counting may only ever lower such a bill."""
|
||||
if rows is None:
|
||||
return None
|
||||
raw = (ep or {}).get("cost") or {}
|
||||
per_row = mk.unit_micro
|
||||
if raw.get("currency") == "credit":
|
||||
try:
|
||||
credits = credits_per_row if credits_per_row is not None else Decimal(str(raw.get("value", 1)))
|
||||
per_row = int(credits * mk.unit_micro)
|
||||
except (InvalidOperation, ValueError):
|
||||
return None
|
||||
return min(rows * per_row, mk.estimate_micro)
|
||||
|
||||
|
||||
def _tavily_requested_result_limit(mk: MarketplaceCall) -> int:
|
||||
"""The request-bound maximum frozen before relay; malformed evidence keeps the 20-page cap."""
|
||||
request = mk.request_data.get("body") if isinstance(mk.request_data, dict) else None
|
||||
@@ -335,6 +409,9 @@ def _observed_cost_micro(mk: MarketplaceCall, body: bytes, headers=None) -> int
|
||||
- fiber-ai: REPORTED in credits, `chargeInfo.creditsCharged` on every envelope, honoured
|
||||
for `method: charged-now` only (a poll repeats its job's charge). Error bodies carry no
|
||||
`chargeInfo`, which is what keeps a 400/404 on a `per_call` profile fetch unbilled.
|
||||
- companyenrich / icypeas bulk / serpstat / thecompaniesapi search / findymail employees:
|
||||
DERIVED by counting the rows the vendor bills for, priced at the row's credits and capped
|
||||
at the hold (`_rows_billed_micro`): an empty answer never costs the requested page.
|
||||
Everyone else settles at the estimate. This is the same signal the catalog's `observed_cost`
|
||||
harvests, which is what lets phase 5's drift detector compare the two numbers directly."""
|
||||
provider = mk.provider
|
||||
@@ -383,6 +460,34 @@ def _observed_cost_micro(mk: MarketplaceCall, body: bytes, headers=None) -> int
|
||||
if isinstance(doc, list) and mk.unit_micro > 0:
|
||||
return sum(item is not None for item in doc) * mk.unit_micro
|
||||
return None
|
||||
if provider == "companyenrich" and mk.cost_type == "per_result" and mk.unit_micro > 0:
|
||||
return _rows_billed_micro(mk, ep, _companyenrich_record_count(mk.endpoint_id, doc))
|
||||
if provider == "icypeas" and mk.cost_type == "per_result" and mk.unit_micro > 0:
|
||||
body = mk.request_data.get("body") if isinstance(mk.request_data, dict) else None
|
||||
# The scrape row carries the dearer profile rate; a company batch is 0.5 credit a hit.
|
||||
company = mk.endpoint_id == "icypeas.scrape.bulk" and isinstance(body, dict) \
|
||||
and body.get("type") == "company"
|
||||
return _rows_billed_micro(mk, ep, _icypeas_bulk_found_count(mk.endpoint_id, doc),
|
||||
Decimal("0.5") if company else None)
|
||||
if provider == "serpstat" and mk.cost_type == "per_result" and mk.unit_micro > 0:
|
||||
return _rows_billed_micro(mk, ep, _serpstat_result_count(doc))
|
||||
if provider == "thecompaniesapi":
|
||||
# `simplified=true` returns a reduced record for zero credits on the endpoints that declare
|
||||
# it (catalog notes); otherwise the company search bills one credit per company RETURNED,
|
||||
# while the reserve is the requested `size`.
|
||||
query_params = (mk.request_data.get("queryParams") or {}) if isinstance(mk.request_data, dict) else {}
|
||||
declared = ((ep or {}).get("input") or {}).get("queryParams") or {}
|
||||
if "simplified" in declared and str(query_params.get("simplified")).lower() == "true":
|
||||
return 0
|
||||
if mk.endpoint_id == "thecompaniesapi.companies.search" and mk.cost_type == "per_result" \
|
||||
and mk.unit_micro > 0 and isinstance(doc, dict) and isinstance(doc.get("companies"), list):
|
||||
return _rows_billed_micro(mk, ep, len(doc["companies"]))
|
||||
if provider == "findymail" and mk.endpoint_id == "findymail.search.employees":
|
||||
# One finder credit per contact RETURNED, and the body is the bare list: an empty `[]` is a
|
||||
# free miss, where the estimate billed the hold.
|
||||
if isinstance(doc, list) and mk.cost_type == "per_result" and mk.unit_micro > 0:
|
||||
return _rows_billed_micro(mk, ep, sum(item is not None for item in doc))
|
||||
return None
|
||||
if not isinstance(doc, dict):
|
||||
return 0 if provider == "contactout" else None
|
||||
reported = (ep.get("cost") or {}).get("reported_charge") if ep else None
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
"""Age out failed-call evidence (`callrecord.error_request` / `error_response`) past retention.
|
||||
|
||||
Run by the `treg-worker admin purge-evidence` cron, never on a request: `GET /admin/errors` is a
|
||||
read, and an admin reading errors during an incident must not blank evidence platform-wide. Until
|
||||
the cron has run, the view itself withholds evidence older than the window (routers/admin.py), so
|
||||
a late or missing schedule delays the UPDATE but never extends what a reader can see.
|
||||
|
||||
An UPDATE, not a DELETE: `callrecord` is the audit trail and the rest of the row must survive. The
|
||||
sentinel rather than NULL keeps "captured, then aged out" distinguishable from "never captured" —
|
||||
without it an old failure and a successful call look identical.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import timedelta
|
||||
|
||||
from sqlalchemy import func, or_, select, update
|
||||
|
||||
from ..infra.db import session_maker
|
||||
from ..models import CallRecord
|
||||
from ..timeutil import utcnow_naive
|
||||
|
||||
ERROR_EVIDENCE_TTL_DAYS = 14
|
||||
ERROR_EVIDENCE_EXPIRED = "<expired>"
|
||||
|
||||
|
||||
def cutoff():
|
||||
return utcnow_naive() - timedelta(days=ERROR_EVIDENCE_TTL_DAYS)
|
||||
|
||||
|
||||
async def purge(batch_size: int = 5000, session_factory=session_maker) -> dict:
|
||||
"""Blank expired evidence `batch_size` rows per transaction → {purged, batches, error}.
|
||||
|
||||
Each batch is its own short transaction on a bounded id set, so no run holds a lock over the
|
||||
whole backlog, and a failure keeps the batches already committed (the next run resumes, because
|
||||
purged rows no longer match). Overlapping runs are harmless: the UPDATE is idempotent."""
|
||||
if batch_size < 1:
|
||||
raise ValueError("batch_size must be >= 1")
|
||||
limit = cutoff()
|
||||
# `coalesce`, not a bare `!=`: SQL three-valued logic makes `error_response != '<expired>'`
|
||||
# UNKNOWN when that column is NULL, so a row carrying request-only evidence would never age
|
||||
# out — excluded by the very predicate meant only to skip rows already purged.
|
||||
pending = (CallRecord.created_at < limit,
|
||||
or_(CallRecord.error_request.is_not(None), CallRecord.error_response.is_not(None)),
|
||||
or_(func.coalesce(CallRecord.error_request, "") != ERROR_EVIDENCE_EXPIRED,
|
||||
func.coalesce(CallRecord.error_response, "") != ERROR_EVIDENCE_EXPIRED))
|
||||
purged = batches = 0
|
||||
try:
|
||||
while True:
|
||||
async with session_factory() as db:
|
||||
ids = (select(CallRecord.id).where(*pending)
|
||||
.order_by(CallRecord.id).limit(batch_size).scalar_subquery())
|
||||
result = await db.execute(
|
||||
update(CallRecord).where(CallRecord.id.in_(ids))
|
||||
.values(error_request=ERROR_EVIDENCE_EXPIRED, error_response=ERROR_EVIDENCE_EXPIRED)
|
||||
.execution_options(synchronize_session=False))
|
||||
await db.commit()
|
||||
count = int(result.rowcount or 0)
|
||||
purged += count
|
||||
batches += 1
|
||||
if count < batch_size:
|
||||
break
|
||||
except Exception as exc: # noqa: BLE001 — reported, and the worker exits non-zero on it
|
||||
logging.getLogger("treg").warning("error-evidence purge failed: %s", exc)
|
||||
return {"purged": purged, "batches": batches, "error": str(exc)}
|
||||
return {"purged": purged, "batches": batches, "error": None}
|
||||
@@ -2236,12 +2236,8 @@ adapters:
|
||||
const: {queryParams.page: "1", queryParams.enrich: "false"}
|
||||
out: {people: items, count: totalResults}
|
||||
miss: "items == []"
|
||||
companyenrich.people.search.scroll:
|
||||
accepts: [[company_domain, title], [company_domain]]
|
||||
in: {company_domain: body.domains, title: body.positionQuery}
|
||||
in_expr: {body.domains: "list(company_domain)", body.positionQuery: "list(title)", body.pageSize: limit}
|
||||
out: {people: items, count: totalItems, next_cursor: nextCursor}
|
||||
miss: "items == []"
|
||||
# No adapter for companyenrich.people.search.scroll: its first page is the same query against the
|
||||
# same index as companyenrich.people.search, so routing both billed one empty answer twice.
|
||||
crustdata.people.search:
|
||||
accepts: [[full_name]]
|
||||
in: {full_name: body.filters.value}
|
||||
|
||||
@@ -275,6 +275,8 @@ contracts:
|
||||
location: {type: str, default: null, note: "free-text place ('London, United Kingdom', 'Greater Phoenix') passed through to providers that take one; finer than country, never normalised"}
|
||||
limit: {type: int, default: 10, note: "rows to return — THE price dial: most providers bill per row"}
|
||||
keywords: {type: list, default: null, note: "skills, topics or domain terms the person must match ('microservices', 'match analysis') — the SUBSTANCE of most briefs; a provider with no keyword field answers a looser question and ranks below one that has it"}
|
||||
# A company-blind provider (lusha, wiza: title only) returns the same strangers for every domain.
|
||||
scoping: [company_domain]
|
||||
output:
|
||||
people: {type: list, required: true, note: "the provider's own people rows (name, title, profile url, location…) — shapes differ per provider; raw is the same body"}
|
||||
count: {type: int}
|
||||
|
||||
@@ -420,6 +420,8 @@ endpoints:
|
||||
checked: '2026-08-20'
|
||||
confidence: documented
|
||||
note: "1 credit per found URL; a miss is free"
|
||||
# A miss is HTTP 200 with a status other than FOUND; this is what settles it at zero.
|
||||
expect: {json_path: status, equals: FOUND}
|
||||
verified: '2026-08-20'
|
||||
example_response: examples/icypeas.people.profile.url.json
|
||||
docs_url: https://api-doc.icypeas.com/scrape/user-profile-url-search
|
||||
@@ -450,6 +452,8 @@ endpoints:
|
||||
checked: '2026-08-20'
|
||||
confidence: documented
|
||||
note: "1 credit per found URL; a miss is free"
|
||||
# A miss is HTTP 200 with a status other than FOUND; this is what settles it at zero.
|
||||
expect: {json_path: status, equals: FOUND}
|
||||
verified: '2026-08-20'
|
||||
example_response: examples/icypeas.companies.profile.url.json
|
||||
docs_url: https://api-doc.icypeas.com/scrape/company-profile-url-search
|
||||
|
||||
@@ -196,6 +196,10 @@ endpoints:
|
||||
checked: '2026-09-16'
|
||||
confidence: verified
|
||||
note: One credit for a non-empty 25-result page; free=true deduplicated pages and no-result responses are free.
|
||||
miss:
|
||||
status: 400
|
||||
when: "error_code == 'NO_RESULTS'"
|
||||
means: "no people matched the search filters — a normal outcome (and free), not an API failure (400 with error_code NO_RESULTS; other 400 bodies are real request errors)"
|
||||
docs_url: https://prospeo.io/api-docs/search-person
|
||||
verified: '2026-09-16'
|
||||
example_response: examples/prospeo.people.search.json
|
||||
@@ -226,6 +230,10 @@ endpoints:
|
||||
cost:
|
||||
<<: *search_cost
|
||||
source_url: https://prospeo.io/api-docs/search-company
|
||||
miss:
|
||||
status: 400
|
||||
when: "error_code == 'NO_RESULTS'"
|
||||
means: "no companies matched the search filters — a normal outcome (and free), not an API failure (400 with error_code NO_RESULTS; other 400 bodies are real request errors)"
|
||||
docs_url: https://prospeo.io/api-docs/search-company
|
||||
verified: '2026-09-16'
|
||||
example_response: examples/prospeo.companies.search.json
|
||||
|
||||
+40
-3
@@ -55,7 +55,8 @@ CONFIG_PATH = Path(os.environ["TREG_CONFIG"]).expanduser() if os.environ.get("TR
|
||||
# Per-invocation `--org <slug>` override (stripped from argv in main); overrides the active org.
|
||||
_ORG_OVERRIDE: str | None = None
|
||||
# Global `--json` (stripped in main): human-table commands emit the raw JSON instead — one stable
|
||||
# contract for agents/scripts. Commands that already print JSON are unaffected.
|
||||
# contract for agents/scripts. `call` prints one envelope, `{"result": <body>, "_treg": {...}}`,
|
||||
# and nothing on stderr: the charge line a script merged into stdout used to break its parse.
|
||||
_JSON_OVERRIDE: bool = False
|
||||
|
||||
|
||||
@@ -211,7 +212,8 @@ class _RegistryClient(httpx.Client):
|
||||
retry.headers["x-treg-body-encoding"] = "base64"
|
||||
if "content-type" in request.headers: # preserve JSON so the server still parses it after decode
|
||||
retry.headers["content-type"] = request.headers["content-type"]
|
||||
print(" (edge WAF blocked the request body; retrying base64-encoded)", file=sys.stderr)
|
||||
if not _JSON_OVERRIDE: # `--json` promises a silent stderr
|
||||
print(" (edge WAF blocked the request body; retrying base64-encoded)", file=sys.stderr)
|
||||
return super().send(retry, **kwargs)
|
||||
|
||||
|
||||
@@ -2487,8 +2489,43 @@ def _print_raw_response(response: httpx.Response) -> None:
|
||||
sys.stdout.flush()
|
||||
|
||||
|
||||
def _call_envelope(response: httpx.Response, content_type: str) -> dict:
|
||||
"""`treg --json call`: the provider body under `result` (parsed JSON, text, or base64 for
|
||||
binary) and what the charge and hint lines would have said under `_treg`, in integer micro-USD."""
|
||||
headers = getattr(response, "headers", {}) or {}
|
||||
if not content_type or content_type == "application/json" or content_type.endswith("+json"):
|
||||
try:
|
||||
result = response.json()
|
||||
except ValueError:
|
||||
result = response.text
|
||||
elif content_type.startswith("text/"):
|
||||
result = response.text
|
||||
else:
|
||||
import base64
|
||||
result = {"base64": base64.b64encode(response.content).decode("ascii"), "content_type": content_type}
|
||||
meta: dict = {"http_status": response.status_code}
|
||||
if call_id := headers.get("X-Treg-Call-Id"):
|
||||
meta["call_id"] = call_id
|
||||
if (cost := headers.get("X-Treg-Cost-Micro")) is not None:
|
||||
# An async submission's cost header is a hold pending settlement, not a charge.
|
||||
meta["reserved_micro" if headers.get("X-Treg-Async") else "charged_micro"] = int(cost)
|
||||
if headers.get("X-Treg-Idempotent-Replay"):
|
||||
meta["replay"] = True
|
||||
if headers.get("X-Treg-Async"):
|
||||
meta["async"] = True
|
||||
if kind := headers.get("X-Treg-Hint"):
|
||||
meta["hint"] = kind
|
||||
return {"result": result, "_treg": meta}
|
||||
|
||||
|
||||
def _show_call_response(response: httpx.Response) -> None:
|
||||
content_type = getattr(response, "headers", {}).get("content-type", "").partition(";")[0].strip().lower()
|
||||
if _JSON_OVERRIDE:
|
||||
print(json.dumps(_call_envelope(response, content_type), ensure_ascii=False, separators=(",", ":")))
|
||||
sys.stdout.flush()
|
||||
if response.status_code >= 400:
|
||||
raise SystemExit(1)
|
||||
return
|
||||
if content_type and content_type != "application/json" and not content_type.endswith("+json") \
|
||||
and not content_type.startswith("text/"):
|
||||
sys.stdout.buffer.write(response.content)
|
||||
@@ -5619,7 +5656,7 @@ _GLOBAL_OPTS = [
|
||||
("-h, --help", "Show this help and exit."),
|
||||
("--version", "Print the treg version and exit."),
|
||||
("--org <slug>", "Run any command in that team instead of the active one."),
|
||||
("--json", "Table-rendering commands (org ls, agents ls, catalog, …) print raw JSON instead."),
|
||||
("--json", "Table commands print raw JSON; `call` prints one {result, _treg} line, nothing on stderr."),
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -37,6 +37,10 @@ class Contract:
|
||||
# `treg.<capability>` row is ever generated from it, however many children verify. Used where
|
||||
# the "children" are one provider's price tiers, which are not a choice treg should make.
|
||||
routed: bool = True
|
||||
# Identity keys that SCOPE the answer rather than describe it: an adapter with no place for one
|
||||
# the caller sent answers about someone else entirely (a title-only search asked for the CEO of
|
||||
# one company returns CEOs of any company), so the router drops it instead of ranking it down.
|
||||
scoping: tuple[str, ...] = ()
|
||||
|
||||
@property
|
||||
def required_output(self) -> tuple[str, ...]:
|
||||
@@ -119,6 +123,11 @@ def parse_contracts(doc: dict) -> dict[str, Contract]:
|
||||
for v in c.get("identity") or []:
|
||||
if isinstance(v, dict):
|
||||
types.update({k: str(t) for k, t in v.items()})
|
||||
scoping = c.get("scoping") or []
|
||||
if not isinstance(scoping, list) or any(k not in types for k in scoping):
|
||||
# A typo'd key would silently scope nothing; failing the routing load instead makes
|
||||
# every `treg.*` capability disappear, which test_routing notices at once.
|
||||
raise ValueError(f"contract {cap}: scoping must list identity keys of the contract, got {scoping!r}")
|
||||
out[cap] = Contract(
|
||||
capability=cap, summary=str(c.get("summary") or ""), identity=_variants(c.get("identity")),
|
||||
identity_types=types, derive=dict(c.get("derive") or {}),
|
||||
@@ -127,7 +136,8 @@ def parse_contracts(doc: dict) -> dict[str, Contract]:
|
||||
miss=str(c.get("miss") or ""), idempotent=bool(c.get("idempotent", True)),
|
||||
default_max_cost_usd=(float(c["default_max_cost_usd"]) if c.get("default_max_cost_usd") is not None else None),
|
||||
advice_unverified=str(c.get("advice_unverified") or ""),
|
||||
routed=bool(c.get("routed", True)))
|
||||
routed=bool(c.get("routed", True)),
|
||||
scoping=tuple(str(k) for k in scoping))
|
||||
return out
|
||||
|
||||
|
||||
|
||||
@@ -14,13 +14,26 @@ MAX_ERROR_FALLBACKS = 2
|
||||
MIN_HIT_SAMPLES = 50
|
||||
|
||||
|
||||
def _used_keys(adapter: Adapter) -> set[str]:
|
||||
return set(adapter.in_map) | {n for e in (adapter.in_expr or {}).values() for n in re.findall(r"[A-Za-z_]\w*", e)}
|
||||
|
||||
|
||||
def ignored_filters(adapter: Adapter, contract: Contract, identity: dict[str, Any]) -> tuple[str, ...]:
|
||||
"""Filters the caller supplied that this adapter has no place for — the provider will answer a
|
||||
LOOSER question than the one asked. Pure, and knowable before the call, so ranking can use it."""
|
||||
used = set(adapter.in_map) | {n for e in (adapter.in_expr or {}).values() for n in re.findall(r"[A-Za-z_]\w*", e)}
|
||||
used = _used_keys(adapter)
|
||||
return tuple(k for k in (contract.filters or ()) if identity.get(k) not in (None, "") and k not in used)
|
||||
|
||||
|
||||
def unscoped(adapter: Adapter, contract: Contract, identity: dict[str, Any]) -> tuple[str, ...]:
|
||||
"""The contract's `scoping` keys the caller supplied that this adapter never sends. Unlike an
|
||||
ignored filter this is not a looser answer but a different question, so the candidate is
|
||||
dropped, not ranked down: a `{title}`-only people search asked for `{company_domain, title}`
|
||||
returns the same title-matched strangers for every company, and each one bills as a hit."""
|
||||
used = _used_keys(adapter)
|
||||
return tuple(k for k in contract.scoping if identity.get(k) not in (None, "") and k not in used)
|
||||
|
||||
|
||||
def cost_at(cost_view: dict | None, request: dict | None = None, adapter: Adapter | None = None) -> int | None:
|
||||
"""Micro-USD this request will cost at its requested size (plan §3, bench 08-27): per-result
|
||||
prices × the requested `limit` (default 1 for a lookup); flat per-call/per-success as listed;
|
||||
|
||||
@@ -37,8 +37,8 @@ _is_sqlite = _db_url.startswith("sqlite")
|
||||
#
|
||||
# `admin` is 3: one slot for a staff page left polling, one for a human using the dashboard, one
|
||||
# spare. Any handler that opens a SECOND session eats two, which is how `/admin/errors` ate the
|
||||
# whole pool when this was 2 — see `_purge_expired_error_evidence`, now on `background` and
|
||||
# single-flighted so concurrent readers cannot multiply it.
|
||||
# whole pool when this was 2 — its retention sweep then moved off the request path entirely, to the
|
||||
# `treg-worker admin purge-evidence` cron (application/evidence_retention.py).
|
||||
#
|
||||
# Sizes are PER PROCESS, the web service runs two uvicorn workers, and a rolling deploy runs two
|
||||
# instances, so `per_process × 2 × 2` is what must stay under the database plan's 103 ceiling —
|
||||
@@ -57,7 +57,6 @@ BACKGROUND_CONSUMERS: dict[str, int] = {
|
||||
"archive.prune_worker": 1, # holds one across a whole sweep
|
||||
"archive.refresh_worker": 1,
|
||||
"catalog observation refresh": 1, # singleflight, one task per process
|
||||
"admin evidence sweep": 1, # single-flighted in routers/admin.py
|
||||
"api_keys last used": 1, # throttled best-effort managed-key display metadata
|
||||
}
|
||||
|
||||
|
||||
+14
-56
@@ -2,7 +2,6 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
from datetime import timedelta
|
||||
from decimal import Decimal, InvalidOperation
|
||||
import logging
|
||||
@@ -10,15 +9,16 @@ import logging
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query
|
||||
from fastapi.responses import FileResponse, HTMLResponse
|
||||
from pydantic import BaseModel
|
||||
from sqlalchemy import func, or_, update
|
||||
from sqlalchemy import func, or_
|
||||
from sqlalchemy.orm import defer
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlmodel import select
|
||||
|
||||
from .. import reconcile
|
||||
from ..application import evidence_retention
|
||||
from ..config import get_settings
|
||||
from ..infra import kv
|
||||
from ..infra.db import background_session_maker, get_admin_session
|
||||
from ..infra.db import get_admin_session
|
||||
from ..domain import money
|
||||
from ..models import ArchiveEndpointStat, ArchiveKey, ArchiveSnapshot, Bundle, CallRecord, LedgerEntry, Membership, Org, Referral, Secret, Tool, User
|
||||
from ..timeutil import as_naive as _as_naive
|
||||
@@ -185,8 +185,8 @@ async def admin_calls(
|
||||
"method": c.method, "status": c.status_code, "at": c.created_at.isoformat()} for c in rows]
|
||||
|
||||
|
||||
_ERROR_EVIDENCE_TTL_DAYS = 14
|
||||
_ERROR_EVIDENCE_EXPIRED = "<expired>"
|
||||
_ERROR_EVIDENCE_TTL_DAYS = evidence_retention.ERROR_EVIDENCE_TTL_DAYS
|
||||
_ERROR_EVIDENCE_EXPIRED = evidence_retention.ERROR_EVIDENCE_EXPIRED
|
||||
|
||||
|
||||
@app.get("/admin/errors")
|
||||
@@ -201,13 +201,11 @@ async def admin_errors(
|
||||
Superadmin-only and deliberately not mirrored on `/calls`: the rows hold customers' request
|
||||
content, so v1 keeps them behind the same door as every other cross-tenant view.
|
||||
|
||||
Ageing happens HERE rather than on the request path. There is no scheduler in this app by design
|
||||
(see the comment above `_claim_idempotent`), and the obvious lazy hook — a marker written on the
|
||||
request session — cannot work: `get_admin_session` never commits, so the marker would roll back and the
|
||||
purge would then run on every single failed call. Doing it on this route costs one UPDATE to the
|
||||
person who came to read errors, which is exactly who wants the stale ones gone.
|
||||
Read-only. Ageing is the `treg-worker admin purge-evidence` cron's job
|
||||
(application/evidence_retention.py); until it has run, a row older than the window is shown as
|
||||
expired with no evidence, so a late schedule never widens what this view reveals.
|
||||
"""
|
||||
purged = await _purge_expired_error_evidence()
|
||||
cutoff = evidence_retention.cutoff()
|
||||
since = _utcnow_naive() - timedelta(days=max(1, min(days, 90)))
|
||||
q = (select(CallRecord)
|
||||
.where(CallRecord.created_at >= since,
|
||||
@@ -225,7 +223,6 @@ async def admin_errors(
|
||||
select(Org).where(Org.id.in_({c.org_id for c in rows if c.org_id is not None})))).scalars().all()}
|
||||
return {
|
||||
"since": since.isoformat(), "days": days, "retention_days": _ERROR_EVIDENCE_TTL_DAYS,
|
||||
"expired_rows_purged": purged,
|
||||
"errors": [{
|
||||
"id": c.id, "call_ref": c.call_ref, "at": c.created_at.isoformat(),
|
||||
"org": omap[c.org_id].slug if c.org_id in omap else None,
|
||||
@@ -238,53 +235,14 @@ async def admin_errors(
|
||||
# An aged-out row holds the sentinel, which is a STATE, not content. Returning it as the
|
||||
# request/response would have a reader treat the word `<expired>` as the provider's
|
||||
# answer; `expired` says the same thing without pretending to be evidence.
|
||||
"request": None if c.error_request == _ERROR_EVIDENCE_EXPIRED else c.error_request,
|
||||
"response": None if c.error_response == _ERROR_EVIDENCE_EXPIRED else c.error_response,
|
||||
"expired": c.error_response == _ERROR_EVIDENCE_EXPIRED,
|
||||
} for c in rows],
|
||||
"request": None if expired else c.error_request,
|
||||
"response": None if expired else c.error_response,
|
||||
"expired": expired,
|
||||
} for c in rows
|
||||
for expired in (c.error_response == _ERROR_EVIDENCE_EXPIRED or c.created_at < cutoff,)],
|
||||
}
|
||||
|
||||
|
||||
_purge_lock = asyncio.Lock()
|
||||
|
||||
|
||||
async def _purge_expired_error_evidence() -> int:
|
||||
"""Blank the evidence columns past the retention window; returns how many rows were cleared.
|
||||
|
||||
An UPDATE, not a DELETE: `callrecord` is the audit trail and the rest of the row must survive.
|
||||
The sentinel rather than NULL keeps "captured, then aged out" distinguishable from "never
|
||||
captured" — without it an old failure and a successful call look identical. Runs on its own
|
||||
session because the request's session is not committed for us, and on the BACKGROUND pool
|
||||
rather than admin's: it is a retention sweep nobody is reading, and nesting a second admin
|
||||
session inside an admin request would hold two of that pool's few slots at once.
|
||||
|
||||
Single-flighted: the sweep is idempotent and driven by whoever happens to open the errors page,
|
||||
so N concurrent readers would otherwise run N identical bulk UPDATEs and hold N background
|
||||
slots. One at a time makes it one entry in `db.BACKGROUND_CONSUMERS` instead of `admin`'s size.
|
||||
"""
|
||||
cutoff = _utcnow_naive() - timedelta(days=_ERROR_EVIDENCE_TTL_DAYS)
|
||||
try:
|
||||
async with _purge_lock, background_session_maker() as db:
|
||||
result = await db.execute(
|
||||
update(CallRecord)
|
||||
# `coalesce`, not a bare `!=`: SQL three-valued logic makes `error_response !=
|
||||
# '<expired>'` UNKNOWN when that column is NULL, so a row carrying request-only
|
||||
# evidence would never age out — excluded by the very predicate meant only to skip
|
||||
# rows already purged.
|
||||
.where(CallRecord.created_at < cutoff,
|
||||
or_(CallRecord.error_request.is_not(None),
|
||||
CallRecord.error_response.is_not(None)),
|
||||
or_(func.coalesce(CallRecord.error_request, "") != _ERROR_EVIDENCE_EXPIRED,
|
||||
func.coalesce(CallRecord.error_response, "") != _ERROR_EVIDENCE_EXPIRED))
|
||||
.values(error_request=_ERROR_EVIDENCE_EXPIRED,
|
||||
error_response=_ERROR_EVIDENCE_EXPIRED))
|
||||
await db.commit()
|
||||
return int(result.rowcount or 0)
|
||||
except Exception as exc: # noqa: BLE001 — retention housekeeping must not break the view
|
||||
logging.getLogger("treg").warning("error-evidence purge failed: %s", exc)
|
||||
return 0
|
||||
|
||||
|
||||
@app.get("/admin/health")
|
||||
async def admin_health(_: str = Depends(require_superadmin), db: AsyncSession = Depends(get_admin_session)) -> list[dict]:
|
||||
rows = (await db.execute(select(Secret).where(Secret.health_status != "ok"))).scalars().all()
|
||||
|
||||
@@ -315,7 +315,12 @@ Money notes for agents:
|
||||
routed calls it also bounds the whole waterfall (default $1 there).
|
||||
- A failed call (HTTP 4xx/5xx from the provider) is relayed unchanged and charged nothing; the
|
||||
CLI prints one stderr line with the status, whose answer it is, and the call id — and on a metered
|
||||
success, `treg: charged $<usd> · call id <id>` (stdout is always the exact body). A `/call/`
|
||||
success, `treg: charged $<usd> · call id <id>` (stdout is always the exact body). **A script that
|
||||
parses the output uses `treg --json call …`**: stdout is then one line,
|
||||
`{"result": <body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and stderr is silent,
|
||||
so merging the two streams cannot break the parse (`--await` output is unchanged). Before a loop
|
||||
over many inputs, run a few and check the parsed results: a parse bug discards answers already
|
||||
paid for. A `/call/`
|
||||
may be answered from treg's archive of earlier answers to the exact same question while that
|
||||
answer is still fresh: the body is the provider's verbatim bytes, the response says
|
||||
`X-Treg-Cache: hit` with `X-Treg-Fetched-At` and `X-Treg-Age`. Your team's first call on a
|
||||
|
||||
@@ -88,6 +88,11 @@ Notes:
|
||||
to the balance — they take priority automatically). A 402 with `error: route_max_cost` is
|
||||
different: YOUR `X-Treg-Route-Max-Cost` header refused the call before anything was charged —
|
||||
ask for fewer rows/targets or raise the ceiling.
|
||||
- **Scripting many calls:** use `treg --json call …`. Stdout is one line,
|
||||
`{"result": <provider body>, "_treg": {"http_status", "call_id", "charged_micro"}}`, and nothing
|
||||
goes to stderr, so a script that merges the streams still parses every answer (`--await` output
|
||||
is unchanged). Run a handful and check the parsed results before looping over the whole list: a
|
||||
parse bug throws away answers that were already billed.
|
||||
- The real charge is the response header `X-Treg-Cost-Micro` (micro-USD), with `X-Treg-Call-Id`
|
||||
as the id to quote. On an asynchronous submission that header is the reserved ceiling; the CLI
|
||||
labels it as a reservation, and the terminal task settles the real charge. The catalog `~$/call`
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
treg-worker arena insights [--max-seconds 110] # fold new audit rows into the Arena aggregate
|
||||
treg-worker catalog stats [--max-rows 500000] # fold new audit rows into per-day endpoint stats
|
||||
treg-worker jev xboost [--posts 60] [--min-likes 150] # the /jev launch-radar demo: X posts <24h -> jev
|
||||
treg-worker admin purge-evidence [--batch-size 5000] # blank expired error evidence past 14-day retention
|
||||
|
||||
Not the light `treg` CLI: these need the server extra (DB, platform keys in the env) and make
|
||||
outbound calls to third parties, so they run as Render cron jobs with the server's env — never as
|
||||
@@ -288,6 +289,24 @@ async def _jev_xboost(args) -> int:
|
||||
return 0
|
||||
|
||||
|
||||
async def _admin_purge_evidence(args) -> int:
|
||||
"""Blank failed-call evidence past the 14-day retention window (application/evidence_retention)."""
|
||||
from .application import evidence_retention
|
||||
from .infra.db import verify_db
|
||||
|
||||
await verify_db()
|
||||
result = await evidence_retention.purge(batch_size=args.batch_size)
|
||||
print(json.dumps(result, sort_keys=True))
|
||||
return 1 if result.get("error") else 0
|
||||
|
||||
|
||||
def _positive_int(value: str) -> int:
|
||||
n = int(value)
|
||||
if n < 1:
|
||||
raise argparse.ArgumentTypeError("must be >= 1")
|
||||
return n
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
ap = argparse.ArgumentParser(prog="treg-worker", description=__doc__)
|
||||
sub = ap.add_subparsers(dest="group", required=True)
|
||||
@@ -336,6 +355,12 @@ def main(argv: list[str] | None = None) -> int:
|
||||
xb.add_argument("--posts", type=int, default=60, help="cap on posts judged, by views")
|
||||
xb.add_argument("--min-likes", type=int, default=150)
|
||||
xb.set_defaults(fn=_jev_xboost)
|
||||
admin = sub.add_parser("admin", help="admin maintenance tasks")
|
||||
adminsub = admin.add_subparsers(dest="cmd", required=True)
|
||||
purge = adminsub.add_parser("purge-evidence", help="blank expired error evidence past the 14-day retention window")
|
||||
purge.add_argument("--batch-size", type=_positive_int, default=5000,
|
||||
help="rows to update per transaction (default 5000)")
|
||||
purge.set_defaults(fn=_admin_purge_evidence)
|
||||
args = ap.parse_args(argv)
|
||||
_need_server()
|
||||
return asyncio.run(args.fn(args))
|
||||
|
||||
@@ -1860,7 +1860,7 @@
|
||||
},
|
||||
"/admin/errors": {
|
||||
"get": {
|
||||
"description": "Failed calls with the evidence to explain them — the caller's request and the provider's own\nanswer (see models.CallRecord.error_request).\n\nSuperadmin-only and deliberately not mirrored on `/calls`: the rows hold customers' request\ncontent, so v1 keeps them behind the same door as every other cross-tenant view.\n\nAgeing happens HERE rather than on the request path. There is no scheduler in this app by design\n(see the comment above `_claim_idempotent`), and the obvious lazy hook — a marker written on the\nrequest session — cannot work: `get_admin_session` never commits, so the marker would roll back and the\npurge would then run on every single failed call. Doing it on this route costs one UPDATE to the\nperson who came to read errors, which is exactly who wants the stale ones gone.",
|
||||
"description": "Failed calls with the evidence to explain them — the caller's request and the provider's own\nanswer (see models.CallRecord.error_request).\n\nSuperadmin-only and deliberately not mirrored on `/calls`: the rows hold customers' request\ncontent, so v1 keeps them behind the same door as every other cross-tenant view.\n\nRead-only. Ageing is the `treg-worker admin purge-evidence` cron's job\n(application/evidence_retention.py); until it has run, a row older than the window is shown as\nexpired with no evidence, so a late schedule never widens what this view reveals.",
|
||||
"operationId": "admin_errors_admin_errors_get",
|
||||
"parameters": [
|
||||
{
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
"""`treg --json call`: stdout is one JSON line a script can parse, and stderr stays silent, so a
|
||||
caller that merges the two streams still parses every result. Driven through `main()`."""
|
||||
import json
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from treg import cli
|
||||
|
||||
|
||||
def _serve(monkeypatch, response: httpx.Response):
|
||||
class FakeClient:
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *a):
|
||||
return False
|
||||
|
||||
def request(self, method, url, params=None, content=None, headers=None):
|
||||
return response
|
||||
|
||||
monkeypatch.setattr(cli, "_client", lambda cfg, **k: FakeClient())
|
||||
monkeypatch.setattr(cli, "_load_config", lambda: {"token": "t", "base_url": "http://x"})
|
||||
|
||||
|
||||
def _charged(status=200, **kw):
|
||||
return httpx.Response(status, headers={"X-Treg-Cost-Micro": "98000", "X-Treg-Call-Id": "abc123",
|
||||
"X-Treg-Hint": "review", **kw.pop("headers", {})}, **kw)
|
||||
|
||||
|
||||
def test_json_call_prints_one_envelope_and_nothing_on_stderr(monkeypatch, capsys):
|
||||
_serve(monkeypatch, _charged(json={"items": [{"name": "A"}]}))
|
||||
cli.main(["--json", "call", "treg.people.search", "--data", '{"title":"CEO"}'])
|
||||
out, err = capsys.readouterr()
|
||||
assert err == ""
|
||||
assert out.count("\n") == 1
|
||||
assert json.loads(out) == {"result": {"items": [{"name": "A"}]}, "_treg": {
|
||||
"http_status": 200, "call_id": "abc123", "charged_micro": 98000, "hint": "review"}}
|
||||
|
||||
|
||||
def test_json_call_error_is_still_one_envelope_and_exits_nonzero(monkeypatch, capsys):
|
||||
_serve(monkeypatch, _charged(402, json={"detail": {"error": "insufficient_balance"}}))
|
||||
with pytest.raises(SystemExit) as exc:
|
||||
cli.main(["call", "treg.people.search", "--json"])
|
||||
assert exc.value.code == 1
|
||||
out, err = capsys.readouterr()
|
||||
assert err == "" and json.loads(out)["_treg"]["http_status"] == 402
|
||||
assert json.loads(out)["result"] == {"detail": {"error": "insufficient_balance"}}
|
||||
|
||||
|
||||
def test_json_call_carries_text_and_binary_bodies(monkeypatch, capsys):
|
||||
_serve(monkeypatch, httpx.Response(200, headers={"content-type": "text/csv"}, text="a,b\n1,2\n"))
|
||||
cli.main(["--json", "call", "x.csv"])
|
||||
assert json.loads(capsys.readouterr().out)["result"] == "a,b\n1,2\n"
|
||||
_serve(monkeypatch, httpx.Response(200, headers={"content-type": "image/png"}, content=b"\x89PNG"))
|
||||
cli.main(["--json", "call", "x.png"])
|
||||
assert json.loads(capsys.readouterr().out)["result"] == {"base64": "iVBORw==", "content_type": "image/png"}
|
||||
|
||||
|
||||
def test_async_submission_reports_a_reservation_not_a_charge(monkeypatch, capsys):
|
||||
_serve(monkeypatch, _charged(202, json={"id": "t1"}, headers={"X-Treg-Async": "{}"}))
|
||||
cli.main(["--json", "call", "x.submit"])
|
||||
meta = json.loads(capsys.readouterr().out)["_treg"]
|
||||
assert meta["reserved_micro"] == 98000 and meta["async"] is True and "charged_micro" not in meta
|
||||
|
||||
|
||||
def test_default_call_output_is_unchanged(monkeypatch, capsys):
|
||||
_serve(monkeypatch, _charged(json={"items": []}))
|
||||
cli.main(["call", "treg.people.search"])
|
||||
out, err = capsys.readouterr()
|
||||
assert json.loads(out) == {"items": []}
|
||||
assert "treg: charged $0.098" in err and "abc123" in err
|
||||
@@ -61,9 +61,8 @@ EXPECTED_MAKERS: dict[str, set[str]] = {
|
||||
"worker.py": {API},
|
||||
# Runs only inside `treg-worker catalog stats`; same reasoning as `worker.py`.
|
||||
"application/catalog_stats.py": {API},
|
||||
# Staff pages take their pool through `Depends(get_admin_session)`, not a maker import; the one
|
||||
# maker here is the retention sweep, which is background work and must not nest inside a request.
|
||||
"routers/admin.py": {BACKGROUND},
|
||||
# Runs only inside `treg-worker admin purge-evidence`; same reasoning as `worker.py`.
|
||||
"application/evidence_retention.py": {API},
|
||||
# Off-request writers.
|
||||
"audit.py": {BACKGROUND},
|
||||
"bootstrap.py": {BACKGROUND},
|
||||
@@ -176,7 +175,6 @@ BACKGROUND_SITES = {
|
||||
"archive_bodies.py:_db_fallback": "archive._store/_touch",
|
||||
"archive.py:prune_once": "archive.prune_worker",
|
||||
"archive.py:refresh_once": "archive.refresh_worker",
|
||||
"routers/admin.py:_purge_expired_error_evidence": "admin evidence sweep",
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -28,6 +28,7 @@ from treg.application.call import evidence as call_evidence
|
||||
from treg.application.call import service as call_service
|
||||
from treg.application.call.types import ReservationFailed, UpstreamResponse
|
||||
from treg.routers import admin as admin_routes
|
||||
from treg.application import evidence_retention
|
||||
from treg.routers import call as call_routes
|
||||
from treg.config import get_settings
|
||||
from treg.infra.db import session_maker
|
||||
@@ -509,6 +510,7 @@ async def test_expired_evidence_is_a_state_not_content(clients: AsyncClient, pla
|
||||
row.created_at = row.created_at - timedelta(days=admin_routes._ERROR_EVIDENCE_TTL_DAYS + 1)
|
||||
db.add(row)
|
||||
await db.commit()
|
||||
await evidence_retention.purge()
|
||||
d = (await clients.get("/admin/errors?days=30", headers=ADMIN)).json()
|
||||
aged = [e for e in d["errors"] if e["expired"]]
|
||||
assert aged, "the row is still listed as a failure"
|
||||
@@ -664,9 +666,67 @@ async def test_evidence_ages_out_but_the_audit_row_survives(clients: AsyncClient
|
||||
await db.commit()
|
||||
call_id, status = row.id, row.status_code
|
||||
|
||||
assert (await clients.get("/admin/errors", headers=ADMIN)).json()["expired_rows_purged"] == 1
|
||||
result = await evidence_retention.purge()
|
||||
assert result["purged"] == 1, "retention worker should blank one row"
|
||||
assert result["error"] is None
|
||||
async with session_maker() as db:
|
||||
aged = await db.get(CallRecord, call_id)
|
||||
assert aged.error_response == admin_routes._ERROR_EVIDENCE_EXPIRED, "aged out, not silently NULL"
|
||||
assert aged.status_code == status, "the rest of the audit row is untouched"
|
||||
assert aged.endpoint_id == EP
|
||||
|
||||
|
||||
async def test_admin_errors_is_read_only_and_withholds_aged_evidence(clients: AsyncClient, platform_on,
|
||||
monkeypatch):
|
||||
"""Reading errors changes nothing (it once blanked every aged row platform-wide on each load),
|
||||
and a row past the window shows no evidence even before the purge cron has reached it."""
|
||||
monkeypatch.setattr(call_service, "relay", _fake_relay(400, b'{"error":"stale failure"}'))
|
||||
await clients.get(f"/call/{EP}?aweme_id=bad")
|
||||
from treg import audit
|
||||
await audit.drain()
|
||||
async with session_maker() as db:
|
||||
row = (await db.execute(
|
||||
select(CallRecord).order_by(CallRecord.id.desc()).limit(1))).scalars().first()
|
||||
row.created_at = row.created_at - timedelta(days=admin_routes._ERROR_EVIDENCE_TTL_DAYS + 1)
|
||||
db.add(row)
|
||||
await db.commit()
|
||||
call_id, original = row.id, row.error_response
|
||||
|
||||
d = (await clients.get("/admin/errors?days=30", headers=ADMIN)).json()
|
||||
listed = next(e for e in d["errors"] if e["id"] == call_id)
|
||||
assert listed["expired"] is True and listed["response"] is None and listed["request"] is None
|
||||
async with session_maker() as db:
|
||||
assert (await db.get(CallRecord, call_id)).error_response == original, "a GET wrote nothing"
|
||||
|
||||
|
||||
async def test_purge_worker_blanks_in_bounded_batches(clients: AsyncClient, platform_on, monkeypatch):
|
||||
"""`treg-worker admin purge-evidence` end to end: each transaction touches at most
|
||||
`--batch-size` rows, the run ends, and a failed run exits non-zero instead of looking green."""
|
||||
from treg import audit, worker
|
||||
monkeypatch.setattr(call_service, "relay", _fake_relay(400, b'{"error":"test failure"}'))
|
||||
for _ in range(5):
|
||||
await clients.get(f"/call/{EP}?aweme_id=test")
|
||||
await audit.drain()
|
||||
async with session_maker() as db:
|
||||
rows = (await db.execute(
|
||||
select(CallRecord).order_by(CallRecord.id.desc()).limit(5))).scalars().all()
|
||||
for row in rows:
|
||||
row.created_at = row.created_at - timedelta(days=admin_routes._ERROR_EVIDENCE_TTL_DAYS + 1)
|
||||
db.add(row)
|
||||
await db.commit()
|
||||
ids = [row.id for row in rows]
|
||||
|
||||
assert await evidence_retention.purge(batch_size=2) == {"purged": 5, "batches": 3, "error": None}
|
||||
async with session_maker() as db:
|
||||
assert {(await db.get(CallRecord, i)).error_response for i in ids} == {
|
||||
admin_routes._ERROR_EVIDENCE_EXPIRED}
|
||||
assert await evidence_retention.purge(batch_size=2) == {"purged": 0, "batches": 1, "error": None}
|
||||
with pytest.raises(ValueError):
|
||||
await evidence_retention.purge(batch_size=0)
|
||||
with pytest.raises(SystemExit):
|
||||
worker.main(["admin", "purge-evidence", "--batch-size", "0"])
|
||||
|
||||
def broken():
|
||||
raise RuntimeError("database unavailable")
|
||||
result = await evidence_retention.purge(session_factory=broken)
|
||||
assert result["error"] == "database unavailable" and result["purged"] == 0
|
||||
|
||||
@@ -4473,3 +4473,96 @@ async def test_trestleiq_missing_required_input_never_reaches_upstream(
|
||||
assert (await clients.get(f"/call/{endpoint}")).status_code == 400
|
||||
await clients.post("/secrets", json={"name": "trestleiq", "value": "OWN-TRESTLEIQ"})
|
||||
assert (await clients.get(f"/call/{endpoint}")).status_code == 400
|
||||
|
||||
|
||||
def _priced(endpoint_id: str, query: dict | None = None, body: dict | None = None, **kw):
|
||||
"""A MarketplaceCall priced the way resolve prices it: the reserve and the per-count unit come
|
||||
from `_marketplace_pricing` over the real catalog row, never from a constant."""
|
||||
catalog = catalog_store.load()
|
||||
ep = catalog.by_id[endpoint_id]
|
||||
cv = catalog.cost_view(ep["cost"], ep["provider"])
|
||||
raw = json.dumps(body).encode() if body is not None else b""
|
||||
estimate, unit = call_resolution._marketplace_pricing(
|
||||
ep["provider"], endpoint_id, cv, call_resolution.QueryValues(tuple((query or {}).items())), raw)
|
||||
mk = _mk(ep["provider"], endpoint_id=endpoint_id, cost_type=ep["cost"]["type"],
|
||||
unit_micro=unit, estimate_micro=estimate, **kw)
|
||||
return mk, estimate, _usd_to_micro_for_test(cv["usd"])
|
||||
|
||||
|
||||
def _usd_to_micro_for_test(usd) -> int:
|
||||
return int(round(float(usd) * 1_000_000))
|
||||
|
||||
|
||||
@pytest.mark.parametrize(("endpoint_id", "query", "req", "body", "rows"), [
|
||||
# CompanyEnrich: 2 credits per person returned, the 2-credit minimum on an empty page.
|
||||
("companyenrich.people.search", None, {"pageSize": 10}, b'{"items":[]}', 1),
|
||||
("companyenrich.people.search", None, {"pageSize": 10}, b'{"items":[{},{},{}]}', 3),
|
||||
("companyenrich.people.search", None, {"pageSize": 10}, b'{"totalItems":0}', None),
|
||||
# Icypeas bulk: only FOUND rows bill, at the row's credits (10 per reverse-email hit).
|
||||
("icypeas.people.identity.resolve.bulk", None, {"data": [["a@x.io"], ["b@x.io"], ["c@x.io"]]},
|
||||
b'{"data":[{"status":"FOUND"},{"status":"NOT_FOUND"},{"status":"FOUND"}]}', 2),
|
||||
("icypeas.profile.url.bulk", None, {"data": [["a"], ["b"]]}, b'{"data":[{"status":"NOT_FOUND"}]}', 0),
|
||||
# Serpstat: an error envelope is free, rows bill with a 1-credit minimum, unknown shapes estimate.
|
||||
("serpstat.web.backlinks.list", None, {"params": {"size": 50}},
|
||||
b'{"id":"1","error":{"code":-32600,"message":"Data not found"}}', 0),
|
||||
("serpstat.web.backlinks.list", None, {"params": {"size": 50}}, b'{"id":"1","result":{"data":[{},{}]}}', 2),
|
||||
("serpstat.web.backlinks.list", None, {"params": {"size": 50}}, b'{"id":"1","result":{"data":[]}}', 1),
|
||||
("serpstat.google.domain.overview", None, {"params": {"domains": ["a.com", "b.com"]}},
|
||||
b'{"id":"1","result":{"a.com":{},"b.com":{}}}', None),
|
||||
# TheCompaniesAPI search: one credit per company returned.
|
||||
("thecompaniesapi.companies.search", {"size": "10"}, None, b'{"companies":[]}', 0),
|
||||
("thecompaniesapi.companies.search", {"size": "10"}, None, b'{"companies":[{},{}]}', 2),
|
||||
# Findymail employee search: one credit per contact, never above the hold.
|
||||
("findymail.search.employees", None, {"website": "x.io", "job_titles": ["CEO"], "count": 5}, b'[]', 0),
|
||||
])
|
||||
def test_per_result_search_settles_on_rows_returned_not_rows_requested(endpoint_id, query, req, body, rows):
|
||||
"""Each reserves the requested page; the body says how many rows the vendor billed. The unit
|
||||
comes from the real pricing path, where a credit-priced row's unit is ONE credit."""
|
||||
mk, estimate, per_row = _priced(endpoint_id, query, req)
|
||||
observed = call_settle._observed_cost_micro(mk, body)
|
||||
if rows is None:
|
||||
assert observed is None
|
||||
else:
|
||||
assert observed == min(rows * per_row, estimate), (observed, per_row, estimate)
|
||||
|
||||
|
||||
def test_row_counts_never_bill_above_the_hold():
|
||||
"""A row whose catalog unit names an input entity reserves per thing asked about; counting
|
||||
returned rows may lower that bill, never raise it."""
|
||||
mk, estimate, _ = _priced("findymail.search.employees", None,
|
||||
{"website": "x.io", "job_titles": ["CEO"], "count": 5})
|
||||
assert call_settle._observed_cost_micro(mk, json.dumps([{"name": str(i)} for i in range(5)]).encode()) == estimate
|
||||
mk, estimate, _ = _priced("serpstat.google.domain.ranked_keywords", None,
|
||||
{"params": {"domain": "a.com", "se": "g_us", "size": 1000}})
|
||||
assert call_settle._observed_cost_micro(
|
||||
mk, json.dumps({"id": "1", "result": {"data": [{}] * 1000}}).encode()) == estimate
|
||||
|
||||
|
||||
def test_thecompaniesapi_simplified_is_free_only_where_declared():
|
||||
mk, _, _ = _priced("thecompaniesapi.companies.search", {"size": "10", "simplified": "true"}, None,
|
||||
request_data={"queryParams": {"size": "10", "simplified": "true"}})
|
||||
assert call_settle._observed_cost_micro(mk, b'{"companies":[{},{}]}') == 0
|
||||
ep = catalog_store.load().by_id["thecompaniesapi.companies.email_pattern"]
|
||||
assert "simplified" not in ((ep.get("input") or {}).get("queryParams") or {})
|
||||
other = _mk("thecompaniesapi", endpoint_id=ep["id"], cost_type=ep["cost"]["type"], unit_micro=9_500,
|
||||
request_data={"queryParams": {"simplified": "true"}})
|
||||
assert call_settle._observed_cost_micro(other, b'{"pattern":"{first}"}') != 0
|
||||
|
||||
|
||||
def test_icypeas_profile_url_miss_settles_at_zero():
|
||||
"""No adapter reads these bodies, so the endpoint's `expect` rule is what makes a miss free."""
|
||||
for endpoint in ("icypeas.people.profile.url", "icypeas.companies.profile.url"):
|
||||
mk = _mk("icypeas", endpoint_id=endpoint, cost_type="per_success", unit_micro=3_800)
|
||||
assert call_settle._observed_cost_micro(mk, b'{"success":true,"result":null,"status":"NOT_FOUND"}') == 0
|
||||
assert call_settle._observed_cost_micro(
|
||||
mk, b'{"success":true,"result":"https://www.linkedin.com/in/x","status":"FOUND"}') is None
|
||||
|
||||
|
||||
def test_icypeas_company_scrape_bills_the_company_rate():
|
||||
body = {"type": "company", "data": ["https://www.linkedin.com/company/a", "https://www.linkedin.com/company/b"]}
|
||||
mk, estimate, per_row = _priced("icypeas.scrape.bulk", None, body, request_data={"body": body})
|
||||
found = b'{"data":[{"status":"FOUND"},{"status":"FOUND"}]}'
|
||||
assert call_settle._observed_cost_micro(mk, found) == 2 * mk.unit_micro // 2 # 0.5 credit each
|
||||
mk, _, _ = _priced("icypeas.scrape.bulk", None, {**body, "type": "profile"},
|
||||
request_data={"body": {**body, "type": "profile"}})
|
||||
assert call_settle._observed_cost_micro(mk, found) == 3 * mk.unit_micro # 1.5 credits each
|
||||
|
||||
@@ -2565,3 +2565,30 @@ def test_linkedin_url_lowercases_the_host_so_the_handle_derives():
|
||||
from treg.domain.catalog.routing import paths as P
|
||||
assert P.linkedin_url("LinkedIn.com/in/Patrick") == "https://linkedin.com/in/Patrick"
|
||||
assert P.linkedin_handle(P.linkedin_url("WWW.LinkedIn.com/in/Patrick")) == "Patrick"
|
||||
|
||||
|
||||
async def test_company_blind_search_provider_is_dropped_not_billed(clients: AsyncClient, platform_on, monkeypatch):
|
||||
"""`people.search` declares `scoping: [company_domain]`. A title-only provider asked for
|
||||
`{company_domain, title}` answers with the same strangers for every company and bills them as
|
||||
a hit, so it leaves the plan instead of ranking last — and still serves a title-only search."""
|
||||
for p in ("LUSHA", "COMPANYENRICH"):
|
||||
monkeypatch.setenv(f"TREG_PLATFORM_KEY_{p}", f"PLATFORM-{p}-KEY")
|
||||
monkeypatch.setenv("TREG_PLATFORM_PROVIDERS", "lusha,companyenrich")
|
||||
get_settings.cache_clear()
|
||||
cat = catalog_store.load()
|
||||
for eid in cat.by_id["treg.people.search"]["routed_children"]:
|
||||
if eid not in ("lusha.people.search", "companyenrich.people.search"):
|
||||
monkeypatch.delitem(cat.adapters, eid, raising=False)
|
||||
seen = []
|
||||
monkeypatch.setattr(call_service, "relay", _relay_by_provider({
|
||||
"lusha": [(200, {"results": [{"name": "Some CEO"}], "pagination": {"total": 1}})],
|
||||
"companyenrich": [(200, {"items": [], "totalItems": 0})]}, seen))
|
||||
r = await clients.post("/call/treg.people.search", json={"company_domain": "example.com", "title": "CEO"})
|
||||
assert [s[0] for s in seen] == ["companyenrich"], seen
|
||||
assert r.status_code == 200 and r.json()["_treg"]["outcome"] == "miss", r.text
|
||||
assert any(d["endpoint_id"] == "lusha.people.search" and "company_domain" in d["why"]
|
||||
for d in r.json()["_treg"]["dropped"]), r.json()["_treg"]
|
||||
seen.clear()
|
||||
r = await clients.post("/call/treg.people.search", json={"title": "CEO"})
|
||||
assert [s[0] for s in seen] == ["lusha"] and r.json()["_treg"]["served_by"] == "lusha.people.search", r.text
|
||||
get_settings.cache_clear()
|
||||
|
||||
Reference in New Issue
Block a user