mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
feat(call): bound and settle Octen platform usage
This commit is contained in:
@@ -227,6 +227,7 @@ Regenerate via `scripts/build-map.py`.
|
||||
| `src/treg/application/call/idempotency.py` | architecture/import-boundaries.md, architecture/money.md, architecture/proxy-model.md, interface/api.md |
|
||||
| `src/treg/application/call/intake.py` | architecture/import-boundaries.md, architecture/money.md, architecture/proxy-model.md, interface/api.md |
|
||||
| `src/treg/application/call/invite.py` | architecture/feedback.md |
|
||||
| `src/treg/application/call/octen.py` | architecture/catalog.md, architecture/money.md |
|
||||
| `src/treg/application/call/overflow.py` | architecture/import-boundaries.md, ops/capacity.md |
|
||||
| `src/treg/application/call/reserve.py` | architecture/import-boundaries.md, architecture/money.md, architecture/proxy-model.md, interface/api.md |
|
||||
| `src/treg/application/call/resolve.py` | architecture/import-boundaries.md, architecture/instagram-oauth.md, architecture/money.md, architecture/multi-tenancy.md, architecture/proxy-model.md, interface/api.md |
|
||||
@@ -642,7 +643,7 @@ Regenerate via `scripts/build-map.py`.
|
||||
| `architecture/ads-conversions.md` | `adsconv.py`, `signup.py`, `adtrack.js`, `gtag.js` |
|
||||
| `architecture/archive.md` | `archive.py`, `hunter.yaml`, `results.py`, `0031_archive_result_admission.py`, `test_cache_result_admission.py`, `archive_bodies.py`, `config.py`, `object_store.py`, `0032_archive_body_storage.py`, `test_archive_r2.py`, `fake_object_store.py`, `smoke_archive_r2.py`, `0002_archive_tables.py`, `0003_callrecord_cached.py`, `0004_archivekey_request_shape.py`, `0011_callrecord_archive_link.py`, `service.py`, `settle.py`, `0039_archive_own_key_and_repeat_pricing.py`, `backfill_call_archive_links.py`, `api.py`, `bootstrap.py`, `admin.py`, `asynctasks.py` |
|
||||
| `architecture/auth-secrets.md` | `injectors.py`, `ssrf.py`, `crypto.py`, `oauth.py`, `__init__.py`, `authorization.py`, `oauth_flow.py`, `refresh.py`, `oauth_exchange.py`, `oauth_refresh.py`, `oauth_providers.py`, `session.js`, `keys.js`, `TeamPage.vue`, `health.py`, `connect.py`, `0051_context_dev_tool_host.py`, `connections.py`, `resources.py`, `__init__.py`, `bindings.py`, `bundles.py`, `api_keys.py`, `access.py`, `api_keys.py`, `test_api_keys.py`, `test_oauth_refresh.py`, `test_financialdatasets.py`, `test_key_providers.py`, `config.py` |
|
||||
| `architecture/catalog.md` | `fetchinio.yaml`, `fetchinio.svg`, `fetchinio.linkedin.user.profile.json`, `fetchinio.linkedin.company.profile.json`, `fetchinio.linkedin.user.posts.json`, `fetchinio.linkedin.user.reactions.json`, `fetchinio.linkedin.post.comments.json`, `fetchinio.linkedin.post.reactions.json`, `fetchinio.linkedin.post.engagement.json`, `fishaudio.yaml`, `fishaudio.tts.s2-1-pro.json`, `fishaudio.voices.create.json`, `fishaudio.voices.discover.json`, `provider_resources.py`, `provider_resources.py`, `provider_resources.py`, `tavily.yaml`, `linkup.yaml`, `you.yaml`, `you.web.search.json`, `you.web.contents.json`, `you.svg`, `linkup.svg`, `linkup.web.search.json`, `linkup.web.fetch.json`, `linkup.web.fetch.structured.json`, `linkup.web.answer.json`, `linkup.web.answer.status.json`, `keenable.yaml`, `olostep.yaml`, `spidercloud.yaml`, `spidercloud.web.scrape.json`, `spidercloud.web.crawl.json`, `spidercloud.web.search.json`, `spidercloud.web.links.json`, `spidercloud.web.unblock.json`, `spidercloud.web.screenshot.json`, `tinyfish.yaml`, `tinyfish.web.search.json`, `tinyfish.web.search.news.json`, `tinyfish.web.search.publications.json`, `tinyfish.web.fetch.json`, `tinyfish.web.agent.run.json`, `tinyfish.web.agent.run.get.json`, `tinyfish.web.agent.run.cancel.json`, `test_tinyfish.py`, `exa.yaml`, `anyapi.extended.yaml`, `adyntel.yaml`, `adyntel.meta-ads.library.advertiser.json`, `adyntel.meta-ads.library.search.json`, `adyntel.linkedin.search.ads.company.json`, `adyntel.linkedin.search.ads.keyword.json`, `adyntel.google.ads.transparency.json`, `adyntel.tiktok-ads.library.search.company.json`, `adyntel.google.domain.keywords.overview.json`, `adyntel.svg`, `trestleiq.yaml`, `financialdatasets.yaml`, `test_financialdatasets.py`, `quickenrich.yaml`, `influencersclub.yaml`, `quickenrich.extended.yaml`, `trykitt.yaml`, `contracts.yaml`, `millionverifier.yaml`, `adapters.yaml`, `capabilities.yaml`, `prospeo.yaml`, `test_route_cost_ceiling.py`, `tomba.yaml`, `__init__.py`, `contracts.py`, `paths.py`, `plan.py`, `synthetic.py`, `async_bridge.py`, `route.py`, `wiza.yaml`, `wiza.people.email.find.json`, `wiza.people.email.find.terminal.json`, `wiza.people.phone.find.json`, `wiza.people.phone.find.terminal.json`, `test_routing.py`, `test_wiza.py`, `catalog-drift.yml`, `catalog_drift.py`, `catalog_ingest.py`, `catalog_validate.py`, `aliases.yaml`, `fx.yaml`, `cloro.yaml`, `aviato.yaml`, `crustdata.yaml`, `google-search-console.yaml`, `google-search-console.extended.yaml`, `google-tag-manager.yaml`, `google-tag-manager.extended.yaml`, `instagram.yaml`, `instagram.extended.yaml`, `justoneapi.extended.yaml`, `minimax.yaml`, `apify.yaml`, `brightdata.yaml`, `companyenrich.yaml`, `oceanio.yaml`, `akta.extended.yaml`, `dataforseo.yaml`, `dataforseo.extended.yaml`, `scrapecreators.yaml`, `scrapecreators.extended.yaml`, `serpapi.yaml`, `serpapi.extended.yaml`, `diffbot.yaml`, `diffbot.extended.yaml`, `tikhub.extended.yaml`, `lusha.extended.yaml`, `openrouter.yaml`, `openrouter.extended.yaml`, `replicate.yaml`, `replicate.extended.yaml`, `reapi.yaml`, `piapi.yaml`, `__init__.py`, `store.py`, `hunter.yaml`, `mcp.py`, `settlement.py`, `stats.py`, `catalog_observations.py`, `catalog_stats.py`, `0038_endpoint_day_stats.py`, `catalog.py`, `test_aigc_pr_b.py`, `test_catalog_api.py`, `test_catalog_validate.py` |
|
||||
| `architecture/catalog.md` | `fetchinio.yaml`, `fetchinio.svg`, `fetchinio.linkedin.user.profile.json`, `fetchinio.linkedin.company.profile.json`, `fetchinio.linkedin.user.posts.json`, `fetchinio.linkedin.user.reactions.json`, `fetchinio.linkedin.post.comments.json`, `fetchinio.linkedin.post.reactions.json`, `fetchinio.linkedin.post.engagement.json`, `fishaudio.yaml`, `fishaudio.tts.s2-1-pro.json`, `fishaudio.voices.create.json`, `fishaudio.voices.discover.json`, `provider_resources.py`, `provider_resources.py`, `provider_resources.py`, `tavily.yaml`, `octen.py`, `linkup.yaml`, `you.yaml`, `you.web.search.json`, `you.web.contents.json`, `you.svg`, `linkup.svg`, `linkup.web.search.json`, `linkup.web.fetch.json`, `linkup.web.fetch.structured.json`, `linkup.web.answer.json`, `linkup.web.answer.status.json`, `keenable.yaml`, `olostep.yaml`, `spidercloud.yaml`, `spidercloud.web.scrape.json`, `spidercloud.web.crawl.json`, `spidercloud.web.search.json`, `spidercloud.web.links.json`, `spidercloud.web.unblock.json`, `spidercloud.web.screenshot.json`, `tinyfish.yaml`, `tinyfish.web.search.json`, `tinyfish.web.search.news.json`, `tinyfish.web.search.publications.json`, `tinyfish.web.fetch.json`, `tinyfish.web.agent.run.json`, `tinyfish.web.agent.run.get.json`, `tinyfish.web.agent.run.cancel.json`, `test_tinyfish.py`, `exa.yaml`, `anyapi.extended.yaml`, `adyntel.yaml`, `adyntel.meta-ads.library.advertiser.json`, `adyntel.meta-ads.library.search.json`, `adyntel.linkedin.search.ads.company.json`, `adyntel.linkedin.search.ads.keyword.json`, `adyntel.google.ads.transparency.json`, `adyntel.tiktok-ads.library.search.company.json`, `adyntel.google.domain.keywords.overview.json`, `adyntel.svg`, `trestleiq.yaml`, `financialdatasets.yaml`, `test_financialdatasets.py`, `quickenrich.yaml`, `influencersclub.yaml`, `quickenrich.extended.yaml`, `trykitt.yaml`, `contracts.yaml`, `millionverifier.yaml`, `adapters.yaml`, `capabilities.yaml`, `prospeo.yaml`, `test_route_cost_ceiling.py`, `tomba.yaml`, `__init__.py`, `contracts.py`, `paths.py`, `plan.py`, `synthetic.py`, `async_bridge.py`, `route.py`, `wiza.yaml`, `wiza.people.email.find.json`, `wiza.people.email.find.terminal.json`, `wiza.people.phone.find.json`, `wiza.people.phone.find.terminal.json`, `test_routing.py`, `test_wiza.py`, `catalog-drift.yml`, `catalog_drift.py`, `catalog_ingest.py`, `catalog_validate.py`, `aliases.yaml`, `fx.yaml`, `cloro.yaml`, `aviato.yaml`, `crustdata.yaml`, `google-search-console.yaml`, `google-search-console.extended.yaml`, `google-tag-manager.yaml`, `google-tag-manager.extended.yaml`, `instagram.yaml`, `instagram.extended.yaml`, `justoneapi.extended.yaml`, `minimax.yaml`, `apify.yaml`, `brightdata.yaml`, `companyenrich.yaml`, `oceanio.yaml`, `akta.extended.yaml`, `dataforseo.yaml`, `dataforseo.extended.yaml`, `scrapecreators.yaml`, `scrapecreators.extended.yaml`, `serpapi.yaml`, `serpapi.extended.yaml`, `diffbot.yaml`, `diffbot.extended.yaml`, `tikhub.extended.yaml`, `lusha.extended.yaml`, `openrouter.yaml`, `openrouter.extended.yaml`, `replicate.yaml`, `replicate.extended.yaml`, `reapi.yaml`, `piapi.yaml`, `__init__.py`, `store.py`, `hunter.yaml`, `mcp.py`, `settlement.py`, `stats.py`, `catalog_observations.py`, `catalog_stats.py`, `0038_endpoint_day_stats.py`, `catalog.py`, `test_aigc_pr_b.py`, `test_catalog_api.py`, `test_catalog_validate.py` |
|
||||
| `architecture/composition.md` | `bootstrap.py`, `bootstrap_handlers.py`, `bootstrap_http.py`, `call_surface.py`, `connect.py`, `mcp_oauth.py`, `session.py`, `admin.py`, `auth.py`, `billing.py`, `call.py`, `connections.py`, `onboard.py`, `orgs.py`, `resources.py`, `referrals.py`, `web.py`, `dump_surface.py`, `test_app_roles.py` |
|
||||
| `architecture/data-model.md` | `0042_pinned_read_scope.py`, `alembic.ini`, `env.py`, `0001_baseline_current_schema.py`, `0002_archive_tables.py`, `0003_callrecord_cached.py`, `0004_archivekey_request_shape.py`, `0005_capacity_policy_snapshot.py`, `0006_overflow_route.py`, `0007_overflow_spend.py`, `0008_org_platform_overflow_disabled.py`, `0009_callrecord_hit.py`, `0017_async_task_record.py`, `0018_async_resource_ownership.py`, `0019_async_poll_failures.py`, `0020_callrecord_created_at_indexes.py`, `0021_ledgerentry_org_created_at_index.py`, `0022_org_spent_today_counter.py`, `0023_callrecord_org_user_created_at_index.py`, `0024_membership_calls_today_counter.py`, `0027_enrich_arena.py`, `0028_arena_insights.py`, `0029_arena_verification_snapshot.py`, `0011_callrecord_archive_link.py`, `0015_idempotentcall_membership_cascade.py`, `0034_managed_api_keys.py`, `0035_default_key_generation.py`, `0036_activity_key_indexes.py`, `0038_endpoint_day_stats.py`, `maintenance.py`, `sitetrack.js`, `models.py`, `0031_archive_result_admission.py`, `0032_archive_body_storage.py`, `0039_archive_own_key_and_repeat_pricing.py`, `0043_provider_resources.py`, `provider_resources.py`, `provider_resources.py`, `0033_signup_promo_eligibility.py`, `0041_searchlog.py`, `timeutil.py`, `db.py`, `referrals.py`, `audit.py`, `evidence_retention.py`, `analytics.py`, `bootstrap_handlers.py`, `ratestore.py`, `auth.py`, `test_postgres_reset.py`, `test_alembic_expand_safety.py`, `test_api_keys.py` |
|
||||
| `architecture/feedback.md` | `feedback_contract.py`, `__init__.py`, `reports.py`, `reviews.py`, `verdicts.py`, `hints.py`, `config.py`, `call.py`, `invite.py`, `kv.py`, `feedback.py`, `feedback.py`, `0025_feedback.py`, `0026_callreview.py`, `0030_feedback_handling.py`, `test_feedback_handling_schema.py`, `feedback.md`, `test_feedback.py`, `test_reviews.py`, `test_endpoint_verdicts.py`, `test_hints.py`, `test_kv.py` |
|
||||
@@ -653,7 +654,7 @@ Regenerate via `scripts/build-map.py`.
|
||||
| `architecture/local-run.md` | `localrun.py`, `egress.py`, `fsjail.py` |
|
||||
| `architecture/mcp-oauth.md` | `auth.py`, `mcp.py`, `health.py`, `mcp_oauth.py`, `session.py`, `access.py`, `api_keys.py`, `api_keys.py`, `auth.py`, `claude-connector.html`, `connect-demo.html`, `CLAUDE-CONNECTOR-SUBMISSION.md`, `test_mcp.py`, `test_mcp_oauth.py`, `test_mcp_directory.py`, `test_marketplace_call.py` |
|
||||
| `architecture/media.md` | `media.py`, `media.py`, `models.py`, `0037_media_hosting.py`, `test_media.py` |
|
||||
| `architecture/money.md` | `tavily.yaml`, `tinyfish.yaml`, `test_tinyfish.py`, `__init__.py`, `settlement.py`, `__init__.py`, `models.py`, `billing.py`, `idempotency.py`, `intake.py`, `resolve.py`, `service.py`, `async_bridge.py`, `route.py`, `reserve.py`, `settle.py`, `tomba.yaml`, `asynctasks.py`, `0017_async_task_record.py`, `0018_async_resource_ownership.py`, `0019_async_poll_failures.py`, `referrals.py`, `budgets.py`, `__init__.py`, `stripe.py`, `reconcile.py`, `referrals.py`, `api.py`, `signup.py`, `promotions.py`, `0033_signup_promo_eligibility.py`, `admin.py`, `billing.py`, `call.py`, `orgs.py`, `referrals.py`, `test_call_architecture.py`, `test_marketplace_call.py`, `test_asynctasks.py` |
|
||||
| `architecture/money.md` | `tavily.yaml`, `tinyfish.yaml`, `test_tinyfish.py`, `__init__.py`, `settlement.py`, `__init__.py`, `models.py`, `billing.py`, `idempotency.py`, `intake.py`, `resolve.py`, `octen.py`, `service.py`, `async_bridge.py`, `route.py`, `reserve.py`, `settle.py`, `tomba.yaml`, `asynctasks.py`, `0017_async_task_record.py`, `0018_async_resource_ownership.py`, `0019_async_poll_failures.py`, `referrals.py`, `budgets.py`, `__init__.py`, `stripe.py`, `reconcile.py`, `referrals.py`, `api.py`, `signup.py`, `promotions.py`, `0033_signup_promo_eligibility.py`, `admin.py`, `billing.py`, `call.py`, `orgs.py`, `referrals.py`, `test_call_architecture.py`, `test_marketplace_call.py`, `test_asynctasks.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_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`, `limits.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` |
|
||||
|
||||
@@ -19,6 +19,7 @@ sources:
|
||||
- src/treg/domain/provider_resources.py
|
||||
- src/treg/routers/provider_resources.py
|
||||
- src/treg/catalog/tavily.yaml
|
||||
- src/treg/application/call/octen.py
|
||||
- src/treg/catalog/linkup.yaml
|
||||
- src/treg/catalog/you.yaml
|
||||
- src/treg/catalog/examples/you.web.search.json
|
||||
@@ -325,6 +326,14 @@ Each endpoint's `tavily_rates` mapping has an exact mode-key contract. Catalog v
|
||||
incomplete, extra, non-finite or non-positive rate, and runtime repeats that check before reserve or
|
||||
relay so catalog drift cannot silently turn a platform call into a free call.
|
||||
|
||||
Octen's `octen_rates` mapping likewise keeps the original published USD unit prices in catalog data.
|
||||
`octen.rates_micro` rejects incomplete or non-micro-USD rates before a shared-key call. The platform
|
||||
request check bounds search result counts, Broad subqueries, News subject results, and Extract URL
|
||||
count without rewriting the caller's body. A team's own key bypasses that check. The call runtime
|
||||
uses the bounded maximum for the hold and this response's `meta.usage` for settlement; the relay
|
||||
does not parse or reshape Octen's answer. Catalog validation requires each route's exact rate keys
|
||||
and checks that the displayed base price matches its rate table.
|
||||
|
||||
Linkup's curated `web.search` and Markdown `web.extract` rows have verified routing adapters;
|
||||
Research remains a direct asynchronous `web.answer` tool. Search's `depth` and `outputType` select
|
||||
one published per-success price. The Fetch route has separate Markdown and structured catalog
|
||||
|
||||
@@ -13,6 +13,7 @@ sources:
|
||||
- src/treg/application/call/idempotency.py
|
||||
- src/treg/application/call/intake.py
|
||||
- src/treg/application/call/resolve.py
|
||||
- src/treg/application/call/octen.py
|
||||
- src/treg/application/call/service.py
|
||||
- src/treg/application/call/async_bridge.py
|
||||
- src/treg/application/call/route.py
|
||||
@@ -524,6 +525,7 @@ Provider-specific calculation stays outside the faithful relay.
|
||||
| Tavily Extract | Reserve the requested URL count (bounded by the documented 20-URL maximum) at 0.2 credit per Basic or 0.4 per Advanced extraction. Settle that fractional allocation for each valid entry in `results`; `failed_results` and grouped `usage.credits` do not charge the caller. A documented empty results list is free; malformed evidence keeps the frozen reserve |
|
||||
| Tavily Map | Platform calls require an explicit integer `limit` from 1 to 20. Reserve that many pages at 0.1 credit each, or 0.2 when the caller supplied nonempty `instructions`; settle valid URL strings in `results` at the frozen per-page unit. Empty results are free, malformed evidence keeps the reserve, and grouped `usage.credits` is ignored |
|
||||
| Tavily Crawl | Platform calls require the same 1-20 limit. Reserve per returned extraction at 0.3 credit (Basic), 0.4 (Basic + instructions), 0.5 (Advanced), or 0.6 (Advanced + instructions), then settle valid extracted entries in `results`. This is a conservative deterministic allocation, not the exact Tavily account charge: the response does not expose every page successfully mapped before extraction. treg absorbs any hidden mapping difference, bounded by the 20-page platform cap. Grouped `usage.credits` is ignored and BYOK remains unmetered |
|
||||
| Octen search and extraction | Reserve the documented maximum from `count`, Broad `max_queries`, optional full-content results, News subjects, or Extract URL count at the advanced rate. Settle Web and News from the call base plus `meta.usage.full_content_extra_count`, Broad from `num_search_queries` plus extras, and Extract from `successful_by_mode`. Missing, malformed, or over-ceiling usage keeps the frozen hold; failed HTTP responses release it. Rates are frozen from `cost.octen_rates` before relay. Own-key calls remain unmetered and have no platform request cap. |
|
||||
| Legacy reported charge | DataForSEO `cost`, ScrapeCreators and Dropleads finder/verifier `credits_charged`, Akta and Dropleads person enrichment `credits_consumed`, Dropleads company `credits.creditsDeducted`, Lusha `billing.creditsCharged`, Exa `costDollars.total`, and Prospeo bulk `total_cost`; credit amounts use the catalog FX rate |
|
||||
| Crustdata, cloro, AI Ark | Read the charge from a response header through `_CREDIT_HEADERS` using the same FX rate. Crustdata `X-Credits-Used` and cloro `X-Credits-Charged` are positive charges; AI Ark `X-Credit` is a negative debit and declares an explicit -1 multiplier. Invalid signs and non-finite values are ignored. cloro omits the header on its free routes and on a failed extraction, neither of which it bills, so an absent header settles at the estimate, not at zero |
|
||||
| cloro reserve | `cost.value` is the full-surface `test_request` price (ChatGPT 9, Google SERP 7); the plain call settles lower from the header (verified live 2026-09-07 at the then-Lite rate: reserve 7,200 µ$, settled 5,600, refunded 1,600; at the Hobby rate 3,600 → 2,800, re-verified 2026-09-14). The top-level `state` body field is a `cost.modifiers` rider (+2 credits) reserved through the same generic path Aviato uses, which is open to any credit-priced provider with a FX rate |
|
||||
|
||||
@@ -299,6 +299,36 @@ def check_tavily_rates(endpoint_id: str, cost: dict, where: str,
|
||||
fail(errors, where, "cost.tavily_rates values must be positive finite numbers")
|
||||
|
||||
|
||||
OCTEN_RATE_KEYS = {
|
||||
"octen.web.search": {"call", "full_content_extra"},
|
||||
"octen.web.search.broad": {"subquery", "full_content_extra"},
|
||||
"octen.web.search.news": {"call", "full_content_extra"},
|
||||
"octen.web.extract": {"standard", "advanced"},
|
||||
}
|
||||
|
||||
|
||||
def check_octen_rates(endpoint_id: str, cost: dict, where: str,
|
||||
errors: list[str]) -> None:
|
||||
"""A platform call needs every Octen component rate before it can reserve."""
|
||||
expected = OCTEN_RATE_KEYS.get(endpoint_id)
|
||||
if expected is None:
|
||||
return
|
||||
rates = cost.get("octen_rates")
|
||||
if not isinstance(rates, dict) or set(rates) != expected:
|
||||
fail(errors, where, "cost.octen_rates must contain exactly " + ", ".join(sorted(expected)))
|
||||
return
|
||||
if any(not _finite_number(value) or value <= 0 or value * 1_000_000 != round(value * 1_000_000)
|
||||
for value in rates.values()):
|
||||
fail(errors, where, "cost.octen_rates values must be positive whole-micro USD rates")
|
||||
return
|
||||
base_key = ("advanced" if endpoint_id == "octen.web.extract" else
|
||||
"subquery" if endpoint_id == "octen.web.search.broad" else "call")
|
||||
if (cost.get("currency") != "USD" or not _finite_number(cost.get("value"))
|
||||
or not _finite_number(cost.get("per")) or cost["per"] <= 0
|
||||
or abs(cost["value"] / cost["per"] - rates[base_key]) > 1e-12):
|
||||
fail(errors, where, "cost.value/per must match the Octen base rate")
|
||||
|
||||
|
||||
def check_platform_auth(ep: dict, where: str, errors: list[str]) -> None:
|
||||
"""Anonymous platform fallback is intentionally narrow: proven public GETs that cost zero."""
|
||||
mode = ep.get("platform_auth")
|
||||
@@ -1232,6 +1262,8 @@ def main(argv: list[str]) -> int:
|
||||
check_cost(cost, where, errors, warnings, inp, service)
|
||||
if service == "tavily":
|
||||
check_tavily_rates(eid, cost, where, errors)
|
||||
if service == "octen":
|
||||
check_octen_rates(eid, cost, where, errors)
|
||||
effective_async = effective_async_descriptor(data.get("async"), ep.get("async"))
|
||||
if effective_async is not None:
|
||||
check_async_descriptor(effective_async, where, str(service), endpoint_index,
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
"""Bound Octen platform holds and price this call's reported usage."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import math
|
||||
|
||||
|
||||
RATE_KEYS = {
|
||||
"octen.web.search": frozenset({"call", "full_content_extra"}),
|
||||
"octen.web.search.broad": frozenset({"subquery", "full_content_extra"}),
|
||||
"octen.web.search.news": frozenset({"call", "full_content_extra"}),
|
||||
"octen.web.extract": frozenset({"standard", "advanced"}),
|
||||
}
|
||||
|
||||
|
||||
def rates_micro(endpoint_id: str, cost: dict) -> dict[str, int]:
|
||||
"""The catalog owns the USD rates; reject incomplete or unsafe declarations."""
|
||||
rates = cost.get("octen_rates")
|
||||
if (not isinstance(rates, dict) or set(rates) != RATE_KEYS[endpoint_id]
|
||||
or any(type(v) not in (int, float) or not math.isfinite(v) or v <= 0
|
||||
or round(v * 1_000_000) != v * 1_000_000 for v in rates.values())):
|
||||
raise ValueError("Octen catalog rates must be positive whole micro-USD values")
|
||||
return {name: round(value * 1_000_000) for name, value in rates.items()}
|
||||
|
||||
|
||||
def _body(body: bytes) -> dict:
|
||||
try:
|
||||
value = json.loads(body)
|
||||
except (ValueError, UnicodeDecodeError):
|
||||
return {}
|
||||
return value if isinstance(value, dict) else {}
|
||||
|
||||
|
||||
def _count(value: object, default: int, maximum: int) -> int:
|
||||
return value if type(value) is int and 1 <= value <= maximum else default
|
||||
|
||||
|
||||
def _extras(options: dict, count: int) -> int:
|
||||
full = options.get("full_content")
|
||||
return max(0, count - 10) if isinstance(full, dict) and full.get("enable") is True else 0
|
||||
|
||||
|
||||
def _news_count(request: dict) -> int:
|
||||
count = _count(request.get("count"), 5, 100)
|
||||
subjects = request.get("subjects")
|
||||
subjects = subjects if isinstance(subjects, dict) else {}
|
||||
if subjects.get("enable", True) is not False:
|
||||
count += _count(subjects.get("count"), 2, 5) * _count(
|
||||
subjects.get("max_sub_news"), 5, 20)
|
||||
return count
|
||||
|
||||
|
||||
def estimate_micro(endpoint_id: str, rates: dict[str, int], body: bytes) -> int:
|
||||
request = _body(body)
|
||||
if endpoint_id == "octen.web.extract":
|
||||
urls = request.get("urls")
|
||||
count = len(urls) if isinstance(urls, list) and 1 <= len(urls) <= 20 else 20
|
||||
return count * rates["advanced"] # auto and uncertain modes can resolve as advanced
|
||||
if endpoint_id == "octen.web.search.broad":
|
||||
queries = _count(request.get("max_queries"), 5, 30)
|
||||
options = request.get("search_options")
|
||||
options = options if isinstance(options, dict) else {}
|
||||
count = _count(options.get("count"), 5, 100)
|
||||
return queries * (rates["subquery"] + _extras(options, count) * rates["full_content_extra"])
|
||||
count = _news_count(request) if endpoint_id == "octen.web.search.news" else _count(
|
||||
request.get("count"), 5, 100)
|
||||
return rates["call"] + _extras(request, count) * rates["full_content_extra"]
|
||||
|
||||
|
||||
def invalid_platform_parameter(endpoint_id: str, body: bytes) -> str | None:
|
||||
"""Only platform calls need a strict spend ceiling; BYOK keeps the upstream's request."""
|
||||
try:
|
||||
request = json.loads(body)
|
||||
except (ValueError, UnicodeDecodeError):
|
||||
return "body"
|
||||
if not isinstance(request, dict):
|
||||
return "body"
|
||||
if endpoint_id == "octen.web.extract":
|
||||
urls = request.get("urls")
|
||||
if not isinstance(urls, list) or not 1 <= len(urls) <= 20 \
|
||||
or any(not isinstance(url, str) or not url for url in urls):
|
||||
return "body.urls"
|
||||
if request.get("mode", "standard") not in ("standard", "advanced", "auto"):
|
||||
return "body.mode"
|
||||
return None
|
||||
if not isinstance(request.get("query"), str) or not request["query"].strip():
|
||||
return "body.query"
|
||||
if endpoint_id == "octen.web.search.broad":
|
||||
if type(request.get("max_queries", 5)) is not int or not 1 <= request.get("max_queries", 5) <= 30:
|
||||
return "body.max_queries"
|
||||
options = request.get("search_options", {})
|
||||
if not isinstance(options, dict):
|
||||
return "body.search_options"
|
||||
prefix = "body.search_options."
|
||||
else:
|
||||
options = request
|
||||
prefix = "body."
|
||||
if type(options.get("count", 5)) is not int or not 1 <= options.get("count", 5) <= 100:
|
||||
return prefix + "count"
|
||||
full = options.get("full_content", {})
|
||||
if not isinstance(full, dict) or ("enable" in full and type(full["enable"]) is not bool):
|
||||
return prefix + "full_content"
|
||||
if endpoint_id == "octen.web.search.news":
|
||||
subjects = request.get("subjects", {})
|
||||
if not isinstance(subjects, dict) or ("enable" in subjects and type(subjects["enable"]) is not bool):
|
||||
return "body.subjects"
|
||||
for name, default, maximum in (("count", 2, 5), ("max_sub_news", 5, 20)):
|
||||
value = subjects.get(name, default)
|
||||
if type(value) is not int or not 1 <= value <= maximum:
|
||||
return "body.subjects." + name
|
||||
return None
|
||||
|
||||
|
||||
def observed_micro(endpoint_id: str, rates: dict[str, int], request: dict,
|
||||
document: object, ceiling: int) -> int | None:
|
||||
"""Missing or malformed usage falls back to the frozen hold, never to zero."""
|
||||
if not isinstance(document, dict):
|
||||
return None
|
||||
code = document.get("code")
|
||||
if type(code) is int and code != 0:
|
||||
return 0
|
||||
if code != 0:
|
||||
return None
|
||||
meta = document.get("meta")
|
||||
usage = meta.get("usage") if isinstance(meta, dict) else None
|
||||
if not isinstance(usage, dict):
|
||||
return None
|
||||
body = request.get("body") if isinstance(request, dict) else None
|
||||
body = body if isinstance(body, dict) else {}
|
||||
|
||||
def integer(name: str) -> int | None:
|
||||
value = usage.get(name)
|
||||
return value if type(value) is int and value >= 0 else None
|
||||
|
||||
if endpoint_id == "octen.web.extract":
|
||||
modes = usage.get("successful_by_mode")
|
||||
if not isinstance(modes, dict):
|
||||
return None
|
||||
standard, advanced = modes.get("standard_urls"), modes.get("advanced_urls")
|
||||
successful = integer("successful_urls")
|
||||
urls = body.get("urls")
|
||||
if any(type(v) is not int or v < 0 for v in (standard, advanced)) \
|
||||
or successful is None or standard + advanced != successful \
|
||||
or not isinstance(urls, list) or successful > len(urls):
|
||||
return None
|
||||
amount = standard * rates["standard"] + advanced * rates["advanced"]
|
||||
else:
|
||||
extra = integer("full_content_extra_count")
|
||||
if extra is None:
|
||||
return None
|
||||
if endpoint_id == "octen.web.search.broad":
|
||||
queries = integer("num_search_queries")
|
||||
maximum = _count(body.get("max_queries"), 5, 30)
|
||||
if queries is None or queries > maximum:
|
||||
return None
|
||||
amount = queries * rates["subquery"] + extra * rates["full_content_extra"]
|
||||
else:
|
||||
amount = rates["call"] + extra * rates["full_content_extra"]
|
||||
return amount if amount <= ceiling else None
|
||||
@@ -580,6 +580,9 @@ _TAVILY_ENDPOINTS = frozenset({
|
||||
"tavily.web.map",
|
||||
"tavily.web.crawl",
|
||||
})
|
||||
_OCTEN_ENDPOINTS = frozenset({
|
||||
"octen.web.search", "octen.web.search.broad", "octen.web.search.news", "octen.web.extract",
|
||||
})
|
||||
_TAVILY_BOUNDED_SITE_ENDPOINTS = frozenset({"tavily.web.map", "tavily.web.crawl"})
|
||||
_TAVILY_PLATFORM_MAX_RESULTS = 20
|
||||
_TAVILY_RATE_KEYS = {
|
||||
@@ -756,6 +759,18 @@ def _marketplace_pricing(
|
||||
"""
|
||||
if not cost:
|
||||
return 0, 0
|
||||
if provider == "octen" and endpoint_id in _OCTEN_ENDPOINTS:
|
||||
from . import octen
|
||||
try:
|
||||
rates = octen.rates_micro(endpoint_id, cost)
|
||||
except ValueError as exc:
|
||||
raise ResolutionFailed(
|
||||
"catalog_price_invalid", status_code=503, detail={
|
||||
"error": "catalog_price_invalid", "endpoint_id": endpoint_id,
|
||||
"message": "Octen pricing is unavailable because its catalog rates are invalid",
|
||||
},
|
||||
) from exc
|
||||
return octen.estimate_micro(endpoint_id, rates, body), 0
|
||||
if provider == "tavily" and endpoint_id in _TAVILY_ENDPOINTS:
|
||||
return _tavily_pricing(endpoint_id, cost, body)
|
||||
if provider == "openmart" and endpoint_id in _OPENMART_METERED_ENDPOINTS:
|
||||
@@ -1527,6 +1542,19 @@ def _enforce_platform_request(ep: dict, body: bytes, headers=None, query=None) -
|
||||
has a singleton enum is the row identity, not caller choice: accepting another value lets a cheap
|
||||
row reserve for an expensive model. Full schema validation remains out of the faithful BYOK path.
|
||||
"""
|
||||
if ep.get("provider") == "octen" and ep.get("id") in _OCTEN_ENDPOINTS:
|
||||
from . import octen
|
||||
invalid = octen.invalid_platform_parameter(ep["id"], body)
|
||||
if invalid:
|
||||
raise ResolutionFailed(
|
||||
"catalog_parameter_invalid", status_code=400, detail={
|
||||
"error": "catalog_parameter_invalid", "endpoint_id": ep["id"],
|
||||
"parameter": invalid,
|
||||
"message": "Octen platform calls require a bounded request; connect your own key "
|
||||
"for the upstream range",
|
||||
},
|
||||
)
|
||||
|
||||
if ep.get("provider") == "openmart" and ep.get("id") in _OPENMART_METERED_ENDPOINTS:
|
||||
requested = _openmart_requested_records(ep["id"], body)
|
||||
parameter = (
|
||||
@@ -2083,6 +2111,13 @@ async def _resolve_marketplace_call(
|
||||
"when": "response", "amount": {"kind": "observed"},
|
||||
"fallback_micro": info_est, "reserve_micro": info_est,
|
||||
}
|
||||
if service == "octen" and ep["id"] in _OCTEN_ENDPOINTS:
|
||||
from . import octen
|
||||
basis = {
|
||||
"when": "response", "amount": {"kind": "observed"},
|
||||
"fallback_micro": info_est, "reserve_micro": info_est,
|
||||
"octen_rates_micro": octen.rates_micro(ep["id"], cv),
|
||||
}
|
||||
common = dict(
|
||||
upstream=upstream, consumed=consumed, endpoint_id=ep["id"], provider=service,
|
||||
params_hash=phash, cost_type=str((ep.get("cost") or {}).get("type") or ""),
|
||||
|
||||
@@ -507,6 +507,13 @@ def _observed_cost_micro(mk: MarketplaceCall, body: bytes, headers=None) -> int
|
||||
# Extract, Map and Crawl intentionally ignore account-grouped usage.credits. Their unit
|
||||
# was frozen from the caller's request mode and only this response's valid results count.
|
||||
return _tavily_cost_micro(mk, doc)
|
||||
if provider == "octen":
|
||||
from . import octen
|
||||
rates = mk.settlement_basis.get("octen_rates_micro")
|
||||
if not isinstance(rates, dict):
|
||||
return None
|
||||
return octen.observed_micro(
|
||||
mk.endpoint_id, rates, mk.request_data, doc, mk.estimate_micro)
|
||||
if provider == "aviato" and mk.endpoint_id == "aviato.people.enrich.bulk":
|
||||
if isinstance(doc, list) and mk.unit_micro > 0:
|
||||
return sum(item is not None for item in doc) * mk.unit_micro
|
||||
|
||||
@@ -684,6 +684,20 @@ def test_tavily_rates_require_complete_positive_finite_endpoint_tables():
|
||||
assert errors
|
||||
|
||||
|
||||
def test_octen_rate_table_matches_displayed_base_and_requires_every_meter():
|
||||
base = {"currency": "USD", "value": 5, "per": 1000,
|
||||
"octen_rates": {"call": 0.005, "full_content_extra": 0.0005}}
|
||||
errors = []
|
||||
validator.check_octen_rates("octen.web.search", base, "test", errors)
|
||||
assert errors == []
|
||||
for cost in (base | {"octen_rates": {"call": 0.005}},
|
||||
base | {"octen_rates": {"call": 0.005, "full_content_extra": -0.0005}},
|
||||
base | {"value": 1}):
|
||||
errors = []
|
||||
validator.check_octen_rates("octen.web.search", cost, "test", errors)
|
||||
assert errors
|
||||
|
||||
|
||||
# ---- ContactOut ----
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,107 @@
|
||||
"""Octen holds bound the upstream bill; settlement uses only this response's usage."""
|
||||
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from treg.application.call import octen
|
||||
from treg.application.call import resolve, settle
|
||||
from treg.application.call.types import ResolutionFailed
|
||||
|
||||
|
||||
RATES = {
|
||||
"octen.web.search": {"call": 5000, "full_content_extra": 500},
|
||||
"octen.web.search.broad": {"subquery": 5000, "full_content_extra": 500},
|
||||
"octen.web.search.news": {"call": 3000, "full_content_extra": 500},
|
||||
"octen.web.extract": {"standard": 1000, "advanced": 2500},
|
||||
}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("endpoint,payload,hold,usage,settled", [
|
||||
("octen.web.search", {"query": "x", "count": 11}, 5000,
|
||||
{"num_search_queries": 1, "full_content_extra_count": 0}, 5000),
|
||||
("octen.web.search", {"query": "x", "count": 11, "full_content": {"enable": True}}, 5500,
|
||||
{"num_search_queries": 1, "full_content_extra_count": 1}, 5500),
|
||||
("octen.web.search.broad", {"query": "x", "max_queries": 4,
|
||||
"search_options": {"count": 12, "full_content": {"enable": True}}},
|
||||
24000, {"num_search_queries": 2, "full_content_extra_count": 1}, 10500),
|
||||
("octen.web.search.news", {"query": "x", "count": 11,
|
||||
"subjects": {"enable": False}, "full_content": {"enable": True}},
|
||||
3500, {"num_search_queries": 1, "num_subject_search_queries": 0,
|
||||
"full_content_extra_count": 1}, 3500),
|
||||
("octen.web.extract", {"urls": ["https://example.com", "https://example.invalid"],
|
||||
"mode": "auto"}, 5000,
|
||||
{"total_urls": 2, "successful_urls": 1,
|
||||
"successful_by_mode": {"standard_urls": 1, "advanced_urls": 0}}, 1000),
|
||||
])
|
||||
def test_hold_and_actual_usage(endpoint, payload, hold, usage, settled):
|
||||
body = json.dumps(payload).encode()
|
||||
assert octen.invalid_platform_parameter(endpoint, body) is None
|
||||
assert octen.estimate_micro(endpoint, RATES[endpoint], body) == hold
|
||||
doc = {"code": 0, "meta": {"usage": usage}}
|
||||
assert octen.observed_micro(endpoint, RATES[endpoint], {"body": payload}, doc, hold) == settled
|
||||
|
||||
|
||||
def test_news_hold_covers_subject_results_when_full_content_is_enabled():
|
||||
request = {"query": "x", "count": 100, "subjects": {"count": 5, "max_sub_news": 20},
|
||||
"full_content": {"enable": True}}
|
||||
assert octen.estimate_micro("octen.web.search.news", RATES["octen.web.search.news"],
|
||||
json.dumps(request).encode()) == 98000
|
||||
|
||||
|
||||
@pytest.mark.parametrize("endpoint,payload,field", [
|
||||
("octen.web.search", {"query": "x", "count": 101}, "body.count"),
|
||||
("octen.web.search", {"query": "x", "full_content": {"enable": "yes"}},
|
||||
"body.full_content"),
|
||||
("octen.web.search.broad", {"query": "x", "max_queries": 31}, "body.max_queries"),
|
||||
("octen.web.search.broad", {"query": "x", "search_options": {"count": 101}},
|
||||
"body.search_options.count"),
|
||||
("octen.web.search.news", {"query": "x", "subjects": {"max_sub_news": 21}},
|
||||
"body.subjects.max_sub_news"),
|
||||
("octen.web.extract", {"urls": ["https://example.com"] * 21}, "body.urls"),
|
||||
])
|
||||
def test_platform_rejects_unbounded_request(endpoint, payload, field):
|
||||
assert octen.invalid_platform_parameter(endpoint, json.dumps(payload).encode()) == field
|
||||
|
||||
|
||||
def test_missing_or_impossible_usage_cannot_settle_as_zero_or_exceed_hold():
|
||||
endpoint = "octen.web.search.broad"
|
||||
request = {"body": {"query": "x", "max_queries": 1}}
|
||||
rates = RATES[endpoint]
|
||||
assert octen.observed_micro(endpoint, rates, request, {"code": 0}, 5000) is None
|
||||
assert octen.observed_micro(endpoint, rates, request, {"code": 400}, 5000) == 0
|
||||
assert octen.observed_micro(endpoint, rates, request, {"code": 0, "meta": {"usage": {
|
||||
"num_search_queries": 2, "full_content_extra_count": 0}}}, 5000) is None
|
||||
assert octen.observed_micro(endpoint, rates, request, {"code": 0, "meta": {"usage": {
|
||||
"num_search_queries": True, "full_content_extra_count": 0}}}, 5000) is None
|
||||
|
||||
|
||||
def test_rate_table_must_be_complete_and_micro_precise():
|
||||
endpoint = "octen.web.extract"
|
||||
assert octen.rates_micro(endpoint, {"octen_rates": {
|
||||
"standard": 0.001, "advanced": 0.0025}}) == RATES[endpoint]
|
||||
for rates in ({"standard": 0.001}, {"standard": 0, "advanced": 0.0025},
|
||||
{"standard": 0.0010001, "advanced": 0.0025}):
|
||||
with pytest.raises(ValueError):
|
||||
octen.rates_micro(endpoint, {"octen_rates": rates})
|
||||
|
||||
|
||||
def test_call_runtime_uses_frozen_octen_rates_and_checks_platform_shape():
|
||||
endpoint = "octen.web.extract"
|
||||
payload = {"urls": ["https://example.com"], "mode": "auto"}
|
||||
body = json.dumps(payload).encode()
|
||||
cost = {"octen_rates": {"standard": 0.001, "advanced": 0.0025}}
|
||||
assert resolve._marketplace_pricing("octen", endpoint, cost, {}, body) == (2500, 0)
|
||||
mk = SimpleNamespace(
|
||||
provider="octen", endpoint_id=endpoint, cost_type="per_success",
|
||||
estimate_micro=2500, request_data={"body": payload},
|
||||
settlement_basis={"octen_rates_micro": RATES[endpoint]},
|
||||
)
|
||||
result = {"code": 0, "meta": {"usage": {"successful_urls": 1,
|
||||
"successful_by_mode": {"standard_urls": 1, "advanced_urls": 0}}}}
|
||||
assert settle._observed_cost_micro(mk, json.dumps(result).encode()) == 1000
|
||||
with pytest.raises(ResolutionFailed) as caught:
|
||||
resolve._enforce_platform_request({"provider": "octen", "id": endpoint},
|
||||
json.dumps({"urls": ["https://example.com"] * 21}).encode())
|
||||
assert caught.value.kind == "catalog_parameter_invalid"
|
||||
Reference in New Issue
Block a user