feat(table): /table/<tool id>, one call answered as rows and columns

For the Google Sheets add-on. /table/ takes exactly the request /call/ takes and runs the same call
on the same road (gates, key injection, hold and settle, audit row, Idempotency-Key); only the
ending differs. /call/ is unchanged: it passes every query parameter to the provider and must
return the provider's answer as is (non-negotiable 4), so this is a sibling route, not an option.

- routers.call.run_call_surface: the body of call_tool, shared by /call/, /catalog/call/ and
  /table/; each passes its own `finish`. /call/ passes _relay_answer (stream unchanged, async
  descriptor, review invitation), so its behavior is the same.
- domain.table: a pure, stdlib-only converter with its own import-linter contract. Routed jobs take
  their columns from the contract (flat: output fields then served_by, a miss is 0 rows; a required
  list: one row per item, people-shaped lists mapped to fixed columns first, the map is data). Hub
  tools take the manifest's output fields. Anything else is flat, list or nested (tables + summary),
  found from the lists of objects in the answer, with one-item wrappers opened. Never raises: an
  answer it cannot shape is raw.
- application.table: reads the answer (8 MiB cap, raw and truncated beyond), asks the provider for
  uncompressed bytes, keeps the call's X-Treg-* headers, maps an upstream non-2xx to a JSON error
  with the same status. No money moves here: the call is settled before the answer returns.
- Behind TREG_TABLE_ENABLED (+ TREG_TABLE_TEAMS / TREG_TABLE_USERS), off by default. With it off
  the route answers a plain 404 that leaves no audit row.

Tests: tests/test_table.py (converter on saved answers; end to end: flag, four shapes, same charge
as /call/, upstream error releases the hold, Idempotency-Key replay costs 0, a hub tool).

Adds docs/context/architecture/table.md; updates composition.md, proxy-model.md, api.md,
import-boundaries.md and the AGENTS.md non-negotiable 4 line.
This commit is contained in:
UncleCode
2026-09-26 14:52:34 +08:00
parent 62035512c9
commit a6dbaee88f
21 changed files with 892 additions and 25 deletions
+8 -3
View File
@@ -234,13 +234,14 @@ Regenerate via `scripts/build-map.py`.
| `src/treg/application/referrals.py` | architecture/money.md, interface/api.md |
| `src/treg/application/search_experiment.py` | architecture/search-experiment.md |
| `src/treg/application/signup.py` | architecture/ads-conversions.md, architecture/money.md, architecture/multi-tenancy.md, interface/api.md |
| `src/treg/application/table.py` | architecture/table.md |
| `src/treg/archive.py` | architecture/archive.md |
| `src/treg/archive_bodies.py` | architecture/archive.md |
| `src/treg/audit.py` | architecture/data-model.md, ops/deploy.md |
| `src/treg/bootstrap.py` | architecture/archive.md, architecture/composition.md, interface/enrich-arena.md |
| `src/treg/bootstrap_handlers.py` | architecture/composition.md, architecture/data-model.md, interface/api.md |
| `src/treg/bootstrap_http.py` | architecture/composition.md, interface/api.md |
| `src/treg/call_surface.py` | architecture/composition.md, architecture/proxy-model.md, interface/api.md |
| `src/treg/call_surface.py` | architecture/composition.md, architecture/proxy-model.md, architecture/table.md, interface/api.md |
| `src/treg/caller_metadata.py` | architecture/multi-tenancy.md, interface/api.md |
| `src/treg/catalog/adapters.yaml` | architecture/catalog.md |
| `src/treg/catalog/adyntel.yaml` | architecture/catalog.md |
@@ -329,7 +330,7 @@ Regenerate via `scripts/build-map.py`.
| `src/treg/cli.py` | architecture/hub.md, architecture/instagram-oauth.md, interface/cli.md, interface/onboarding.md, interface/shell.md |
| `src/treg/cli_analytics.py` | interface/cli.md |
| `src/treg/client_identity.py` | architecture/import-boundaries.md, architecture/proxy-model.md, interface/api.md |
| `src/treg/config.py` | architecture/archive.md, architecture/auth-secrets.md, architecture/feedback.md, architecture/super-admin.md, guides/expanding-a-category.md, ops/deploy.md |
| `src/treg/config.py` | architecture/archive.md, architecture/auth-secrets.md, architecture/feedback.md, architecture/super-admin.md, architecture/table.md, guides/expanding-a-category.md, ops/deploy.md |
| `src/treg/convert.py` | interface/cli.md |
| `src/treg/crypto.py` | architecture/auth-secrets.md |
| `src/treg/domain/__init__.py` | architecture/import-boundaries.md |
@@ -387,6 +388,7 @@ Regenerate via `scripts/build-map.py`.
| `src/treg/domain/money/settlement.py` | architecture/catalog.md, architecture/money.md |
| `src/treg/domain/provider_resources.py` | architecture/catalog.md, architecture/data-model.md, architecture/multi-tenancy.md |
| `src/treg/domain/referrals.py` | architecture/data-model.md, architecture/money.md |
| `src/treg/domain/table/__init__.py` | architecture/table.md |
| `src/treg/domain/tools/__init__.py` | architecture/auth-secrets.md |
| `src/treg/domain/tools/bindings.py` | architecture/auth-secrets.md |
| `src/treg/domain/tools/bundles.py` | architecture/auth-secrets.md, architecture/multi-tenancy.md |
@@ -433,7 +435,7 @@ Regenerate via `scripts/build-map.py`.
| `src/treg/routers/auth.py` | architecture/composition.md, architecture/mcp-oauth.md, architecture/multi-tenancy.md, interface/api.md, interface/enrich-arena.md, interface/onboarding.md |
| `src/treg/routers/auth_helpers.py` | interface/api.md, interface/cli.md |
| `src/treg/routers/billing.py` | architecture/composition.md, architecture/money.md, interface/api.md |
| `src/treg/routers/call.py` | architecture/composition.md, architecture/feedback.md, architecture/instagram-oauth.md, architecture/money.md, architecture/proxy-model.md, interface/api.md |
| `src/treg/routers/call.py` | architecture/composition.md, architecture/feedback.md, architecture/instagram-oauth.md, architecture/money.md, architecture/proxy-model.md, architecture/table.md, interface/api.md |
| `src/treg/routers/catalog.py` | architecture/catalog.md, architecture/hub.md, interface/api.md |
| `src/treg/routers/connections.py` | architecture/auth-secrets.md, architecture/composition.md, guides/expanding-a-category.md, interface/api.md |
| `src/treg/routers/feedback.py` | architecture/feedback.md |
@@ -446,6 +448,7 @@ Regenerate via `scripts/build-map.py`.
| `src/treg/routers/referrals.py` | architecture/composition.md, architecture/money.md, interface/api.md |
| `src/treg/routers/resources.py` | architecture/auth-secrets.md, architecture/composition.md, architecture/multi-tenancy.md, interface/api.md |
| `src/treg/routers/signup_cookies.py` | interface/api.md |
| `src/treg/routers/table.py` | architecture/table.md |
| `src/treg/routers/web.py` | architecture/composition.md, architecture/hub.md, interface/api.md, interface/dashboard.md, interface/landing-sandbox.md, interface/seo.md, interface/skill.md |
| `src/treg/runner.py` | interface/api.md |
| `src/treg/sandbox.py` | interface/landing-sandbox.md |
@@ -586,6 +589,7 @@ Regenerate via `scripts/build-map.py`.
| `tests/test_routing.py` | architecture/catalog.md |
| `tests/test_search_experiment.py` | architecture/search-experiment.md |
| `tests/test_ssrf_public_addresses.py` | architecture/proxy-model.md |
| `tests/test_table.py` | architecture/table.md |
| `tests/test_tag_billing.py` | architecture/proxy-model.md |
| `tests/test_tag_billing_adversarial.py` | architecture/proxy-model.md |
| `tests/test_team_limit.py` | architecture/multi-tenancy.md |
@@ -616,6 +620,7 @@ Regenerate via `scripts/build-map.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` |
| `architecture/super-admin.md` | `api.py`, `admin.py`, `evidence_retention.py`, `access.py`, `config.py` |
| `architecture/table.md` | `__init__.py`, `table.py`, `table.py`, `call.py`, `call_surface.py`, `config.py`, `test_table.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`, `__init__.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` |
+2
View File
@@ -49,6 +49,8 @@ Everything else in this file is guidance; these are the contract, and they win o
Routed endpoints and overflow wrap the child's answer and say so; they never alter it. Responses needing settlement or ownership evidence are buffered by the application
up to 8 MiB; exceeding that limit fails without charging, never returns a successful prefix.
Authorized free final fetches needing no body evidence stream in full.
`/table/` runs the same call through the same road and returns it as rows and columns; it never
changes what `/call/` returns (docs/context/architecture/table.md).
5. Balances change only through money's five entries: grant, topup, reserve, settle, release.
There is deliberately no refund or adjustment entry; an ops correction is a grant.
+1
View File
@@ -35,6 +35,7 @@ covers (frontmatter `sources:`). Regenerate this index with
| [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, evidence_retention.py, access.py, … |
| [The table layer — one call, answered as rows and columns (`/table/`)](architecture/table.md) | built, behind `table_enabled` (TREG_TABLE_ENABLED), off by default | __init__.py, table.py, table.py, call.py, … |
## Interfaces (API · CLI · skill)
+4 -3
View File
@@ -73,8 +73,9 @@ carrying the `build` and `archive_config` fingerprints every server event has (s
`bootstrap_handlers.py` owns the app-wide pool-saturation and HTTP-exception adapters.
`call_surface.split_call_path`
classifies both `/call/` and `/catalog/call/` so those adapters share the same call-id, audit and
idempotency-release contract while retaining `call` versus `catalog_call` ingress attribution. The
classifies `/call/`, `/catalog/call/` and `/table/` so those adapters share the same call-id, audit
and idempotency-release contract while retaining `call`, `catalog_call` and `table` ingress
attribution. The
composition root supplies the call-specific `_stamp_call_exit` callback from `routers/call.py` before registration;
the callback owns call ids, refusal classification, audit fallback, exceptional call telemetry, and
idempotency-label release. After caller identity exists, the pool adapter reports
@@ -111,7 +112,7 @@ architecture test separately pins the dataplane/control startup split and backgr
| Role | HTTP routes and mounts | Background tasks | Startup checks |
|---|---|---|---|
| `all` | The complete surface, including `/run`, static files, `/mcp`, and the flagged `/mcp/v2` | Ads conversion worker when enabled | Read-only DB verify, HTTP client, enabled MCP lifespans |
| `dataplane` | `/call/{rest:path}`, `/catalog/call/{rest:path}`, MCP mounts, and their resource metadata; no `/run`, static files, docs, or OpenAPI | None | Read-only DB verify, HTTP client, enabled MCP lifespans |
| `dataplane` | `/call/{rest:path}`, `/catalog/call/{rest:path}`, `/table/{rest:path}` (see [table](table.md)), MCP mounts, and their resource metadata; no `/run`, static files, docs, or OpenAPI | None | Read-only DB verify, HTTP client, enabled MCP lifespans |
| `control` | Everything except the calling surfaces; includes OAuth issuance, `/run`, and static files | Ads conversion worker when enabled | Read-only DB verify, HTTP client |
No role lifespan writes schema, performs a data backfill, or provisions the local single user. The explicit
@@ -139,6 +139,10 @@ awaiter and the settlement worker can never disagree about what "done" means. Tw
light: the module is on `test_import_lightness`'s list, and an import-linter contract forbids it
every server root (`treg.models`, `treg.infra`, `treg.config`, SQLModel, pydantic, yaml, httpx).
The table domain (`treg.domain.table`, [table](table.md)) is the same kind of leaf: a parsed answer
in, rows and columns out. Its contract forbids every outer layer, every sibling domain, and the web,
database and settings libraries, so the converter stays testable on saved answers alone.
The capacity domain (`treg.domain.capacity`, plan step B) is a leaf like identity: it cannot import
`treg.api`, `treg.routers`, `treg.application`, `treg.bootstrap`, `treg.audit`, FastAPI or Starlette.
It reads config and writes only its own tables and ratestore keys, from worker-profile commands
+3 -1
View File
@@ -184,7 +184,9 @@ pool occupancy at the network boundary and covers concurrent calls. Pool sizing,
the separate API/admin/background pools are specified in [deploy](../ops/deploy.md).
## Tool resolution (`application.call.resolve`)
`* /call/{rest:path}` → `routers.call.call_tool()` → `application.call.service.execute_call()`
`* /call/{rest:path}` → `routers.call.call_tool()` → `routers.call.run_call_surface()` (shared with
`/catalog/call/` and `/table/`, which differ only in the `finish` that turns the answer into the
response: see [table](table.md)) → `application.call.service.execute_call()`
→ `resolve_call_target(...)` returns a framework-neutral
`ResolvedTarget(tool, upstream)`. Each resolution use case owns and closes its read session.
+99
View File
@@ -0,0 +1,99 @@
---
title: The table layer — one call, answered as rows and columns (`/table/`)
status: built, behind `table_enabled` (TREG_TABLE_ENABLED), off by default
sources:
- src/treg/domain/table/__init__.py
- src/treg/application/table.py
- src/treg/routers/table.py
- src/treg/routers/call.py
- src/treg/call_surface.py
- src/treg/config.py
- tests/test_table.py
related:
- architecture/proxy-model.md
- architecture/hub.md
- architecture/catalog.md
- interface/api.md
---
# The table layer
`* /table/{rest:path}` takes exactly the request `/call/<tool id>` takes (method, query, body, key,
`Idempotency-Key`) and answers the same call as `{shape, columns, rows, ...}`, for clients that want
a grid: the Google Sheets add-on first. It is a **sibling** of `/call/`, not an option on it:
`/call/` passes every query parameter through to the provider (a `?format=` would reach it), and it
must return the provider's answer unchanged (AGENTS.md, non-negotiable 4).
## One call road, two endings
`routers.call.run_call_surface(rest, request, caller, prefix=, finish=, headers=)` is the whole call
road for every call surface: the caller's identity stashed for the refusal fallback, the raw path
rebuilt from `raw_path` after `prefix`, the `CallInput`, `create_call_context`, `execute_call`, and
the same bookkeeping on every exit (audit row, idempotency-claim release, `X-Treg-Call-Id`).
`call_tool` passes `_relay_answer` (stream the answer unchanged, attach the async descriptor and the
review invitation); `table_tool` passes `_table_answer`. `finish` runs inside the same `try`, so a
fault while converting is audited and released like any other.
`table_tool` passes `headers=application.table.plain_headers(...)`: the caller's headers with
`Accept-Encoding: identity`, because the table must parse the body (the hub runner asks its steps the
same way). A provider that answers gzip or deflate anyway is decoded (`_decoded`); any other encoding
reads as a `raw` table.
`/table/` is in `call_surface._CALL_SURFACES` (surface label `table`), so its refusals get
`X-Treg-Error: 1` and the audit fallback, and in `bootstrap._DATAPLANE_ROUTE_KEYS`. The demo lockdown
(`domain.identity.access`) allows only `/call/` for non-read methods, so `/table/` stays out of demo
orgs.
## Money
Nothing here holds, charges or refunds. `execute_call` returns with the call already settled (its
cost is on the answer as `X-Treg-Cost-Micro`), so a conversion that fails can only make the view
`raw`: it can never charge twice or leave a hold open. The table answer carries the call's own
`X-Treg-*` headers (call id, cost, served-by, `X-Treg-Idempotent-Replay`), and the same
`Idempotency-Key` returns the same table for nothing. `application.table.read_answer` reads at most
`MAX_TABLE_BYTES` (8 MiB); a bigger answer is `raw`, `truncated: true`, cut to 64 KB.
An upstream non-2xx keeps its status and answers `{error: "upstream_error", upstream_status,
body_excerpt}` (the first 2 KB). A treg refusal (`CallFailure`) is raised exactly as on `/call/`.
## The flag
`application.table.enabled_for(slug, email)`: `table_enabled`, then `table_teams` / `table_users`
(either list lets a caller in; both empty means every team), the same pattern as the hub. A caller
outside answers a plain 404 **returned, not raised**: a raised 404 on a call surface is stamped and
audited as a refusal, and with the flag off the route must leave no row. `llms.txt` and `skill.md`
do not mention `/table/` while the flag is off.
## The converter (`domain.table`)
Pure and stdlib-only (import-linter contract "Table domain is a pure stdlib leaf"). `to_table(body,
contract_output=, list_field=, hub_fields=)` never raises: anything it cannot shape is `raw`.
`application.table.table_answer` chooses the source of columns (`column_source`):
- **contract**: a routed job (`catalog.by_id[id].kind == "routed"`, contract from
`catalog.contracts[capability]`). A flat contract: its `output` fields in contract order, then
`served_by` (from `_treg.served_by`); one row, none when `_treg.outcome == "miss"`. A contract with
a required `list` field (`people`, `companies`, `results`...): shape `list`, one row per item.
Items are the provider's own objects, so `LIST_MAPS` (data, keyed by the list field; `people` today)
maps common names to fixed columns first: `first_name`, `last_name`, `title`, `company`,
`linkedin_url`, `location`, first matching path wins (`organization.name` is a path); every other
field follows under the provider's own name, minus the paths a mapped column used.
- **hub**: an id not in the catalog that `hub.tool_for` resolves for this caller (one short read,
after the answer is read: non-negotiable 3). Columns are the manifest's `output.fields`, or its
`output` keys, in manifest order; one row.
- **generated**: anything else. `_find_tables` collects every list of objects, searching inside a
one-object list whose key is in `ENVELOPE_KEYS` (`summary[0]`, `tasks[0].result[0]`) and never
inside a table's own items. None: shape `flat`, one row of `flatten` (dotted keys, `MAX_DEPTH` 3).
One: shape `list` (with its `path`). Two or more: shape `nested`, with `tables` (name, path,
row_count, columns, rows) and `summary` (`field`, `value` rows of the tables' common parent; a
list of objects is one cell of each item's first value). `columns`/`rows` are the summary, so a
client that reads only those still gets a grid.
Cells (`cell`): a scalar as it is; a list of scalars joined with `", "`; an object or a list holding
objects as compact JSON. A field whose name starts with `_` (`_treg`, a provider's `_note`) is never
a column; a routed answer's `_treg` goes into the reply's `_treg` with `call_id` and `cost_micro`.
**Trap:** the saved example answers (`src/treg/catalog/examples/`) end every long list with the
string `"… N more item(s) truncated"` (written by `scripts/catalog_verify.py`). A list with that
marker is not "all objects", so the converter skips it (`_TRUNCATION_MARK`); real answers never
carry it.
+1 -1
View File
@@ -194,7 +194,7 @@ same platform key, `retry=1` when a burst-429 with a short `retry-after` was re-
## `X-Treg-Error` - whose refusal is this?
`bootstrap_handlers._mark_treg_own_errors` tags treg's **own**
refusals on `/call/` and `/catalog/call/` paths with `X-Treg-Error: 1`, then answers exactly as
refusals on `/call/`, `/catalog/call/` and `/table/` paths with `X-Treg-Error: 1`, then answers exactly as
before - the status and body
are untouched, and a client that ignores the header sees what it always saw. Without it a caller cannot
tell treg's 404 ("no tool registered for that host") from the vendor's own 404: both are a status and
+8
View File
@@ -302,6 +302,14 @@ forbidden_modules = ["treg.api", "treg.routers", "treg.application", "treg.boots
allow_indirect_imports = true
as_packages = true
[[tool.importlinter.contracts]]
name = "Table domain is a pure stdlib leaf: an answer in, rows and columns out"
type = "forbidden"
source_modules = ["treg.domain.table"]
forbidden_modules = ["treg.api", "treg.routers", "treg.application", "treg.bootstrap", "treg.audit", "treg.models", "treg.infra", "treg.config", "treg.domain.identity", "treg.domain.governance", "treg.domain.connections", "treg.domain.tools", "treg.domain.catalog", "treg.domain.capacity", "treg.domain.money", "treg.domain.asynctasks", "treg.domain.hub", "fastapi", "starlette", "sqlmodel", "sqlalchemy", "pydantic", "pydantic_settings", "yaml", "httpx"]
allow_indirect_imports = true
as_packages = true
[[tool.importlinter.contracts]]
name = "Capacity domain does not depend on outer layers"
type = "forbidden"
+2
View File
@@ -65,6 +65,7 @@ from .routers import orgs as org_routes
from .routers import provider_resources as provider_resource_routes
from .routers import referrals as referral_routes
from .routers import resources as resources_routes
from .routers import table as table_routes
from .routers import web as web_routes
from .routers.auth import _client_ip
from .routers.auth_helpers import _same_origin
@@ -923,6 +924,7 @@ router.routes.extend(admin_routes.reports_router.routes)
# ---- the proxy: call a tool without holding its credential; tier-4 metering ----------------
router.routes.extend(call_routes.router.routes)
router.routes.extend(table_routes.router.routes)
router.routes.extend(arena_routes.router.routes)
+162
View File
@@ -0,0 +1,162 @@
"""`/table/<tool id>`: one ordinary call, its answer read in full and returned as rows and columns.
The call itself runs on the call road unchanged (`routers.call.run_call_surface`): the same gates,
the same key injection, the same hold and settle, the same audit row, the same `Idempotency-Key`.
By the time `table_answer` runs, the call is finished and its money is final (the cost is already on
the answer as `X-Treg-Cost-Micro`), so nothing here can charge, refund or hold anything. A table is
a view of an answer that was bought: when the answer cannot be shaped, the view is `raw`, never an
error. `/call/` is not involved and never changes (AGENTS.md, non-negotiable 4).
"""
from __future__ import annotations
import gzip
import json
import zlib
from typing import Any
from ..domain import table as table_domain
from ..domain.catalog import store as catalog_store
from .call.types import CallContext, UpstreamResponse
MAX_TABLE_BYTES = 8 * 1024 * 1024 # the most of one answer /table/ reads (non-negotiable 4's cap)
BODY_EXCERPT_BYTES = 2048 # the part of an upstream error body the table error carries
def enabled_for(org_slug: str | None, email: str | None) -> bool:
"""`/table/` exists for this caller: the flag, then the team and person lists (either lets the
caller in; both empty means every team). A caller outside them gets the 404 the flag off gives."""
from ..config import get_settings
s = get_settings()
if not s.table_enabled:
return False
if not (s.table_team_set or s.table_user_set):
return True
return ((org_slug or "").lower() in s.table_team_set
or (email or "").strip().lower() in s.table_user_set)
def plain_headers(raw_headers: tuple[tuple[bytes, bytes], ...]) -> list[tuple[bytes, bytes]]:
"""The caller's headers for the upstream call, asking for uncompressed bytes (a table must parse
the body), the way the hub runner asks for its steps."""
kept = [(k, v) for k, v in raw_headers if k.lower() != b"accept-encoding"]
return kept + [(b"accept-encoding", b"identity")]
async def read_answer(upstream: UpstreamResponse) -> tuple[bytes, bool]:
"""The whole answer, at most MAX_TABLE_BYTES; (bytes, truncated). The stream is closed either way."""
chunks: list[bytes] = []
size = 0
truncated = False
try:
async for chunk in upstream.body_stream:
if size + len(chunk) > MAX_TABLE_BYTES:
chunks.append(chunk[:MAX_TABLE_BYTES - size])
truncated = True
break
chunks.append(chunk)
size += len(chunk)
finally:
await upstream.close()
return b"".join(chunks), truncated
def _header(upstream: UpstreamResponse, name: str) -> str | None:
wanted = name.lower().encode("latin-1")
for k, v in upstream.raw_headers:
if k.lower() == wanted:
return v.decode("latin-1")
return None
def _decoded(body: bytes, encoding: str | None) -> bytes:
"""A provider that ignored `Accept-Encoding: identity`: gzip and deflate are undone here; any other
encoding stays as it is and reads as a `raw` table."""
enc = (encoding or "").strip().lower()
try:
if enc == "gzip":
return gzip.decompress(body)
if enc == "deflate":
return zlib.decompress(body)
except (OSError, zlib.error):
return body
return body
def _tool_id(context: CallContext, rest: str) -> str:
marketplace = context.marketplace
if marketplace is not None and marketplace.endpoint_id:
return marketplace.endpoint_id
return rest.split("?", 1)[0]
async def _hub_fields(context: CallContext, rest: str) -> list[str] | None:
"""The output fields of the hub tool this call ran, in manifest order, or None when the id is not
a hub tool. One short read, after the answer is fully read (non-negotiable 3)."""
from . import hub as hub_app
from ..infra.db import session_maker
tool_ref = rest.split("?", 1)[0]
if not hub_app.is_hub_id_shape(tool_ref):
return None
caller = context.input.caller
async with session_maker() as db:
row = await hub_app.tool_for(db, tool_ref, caller_org_id=caller.org_id,
caller_slug=caller.org.slug, caller_email=caller.email)
if row is None:
return None
output = (row.manifest or {}).get("output") or {}
if isinstance(output, dict) and isinstance(output.get("fields"), list):
return [str(f) for f in output["fields"]]
return [str(k) for k in output] if isinstance(output, dict) else []
def _contract(tool_id: str) -> tuple[list[str], str | None] | None:
"""(output fields in contract order, the required list field or None) for a routed job."""
cat = catalog_store.load()
ep = cat.by_id.get(tool_id)
if not ep or ep.get("kind") != "routed":
return None
contract = cat.contracts.get(ep.get("capability") or "")
if contract is None:
return None
fields = list(contract.output)
list_field = next((f for f, spec in contract.output.items()
if (spec or {}).get("type") == "list" and (spec or {}).get("required")), None)
return fields, list_field
async def table_answer(context: CallContext, upstream: UpstreamResponse, rest: str) -> tuple[int, dict[str, Any], list[tuple[bytes, bytes]]]:
"""(status, JSON body, headers) for one finished call. The headers are the call's own `X-Treg-*`
headers (call id, cost, served-by, idempotent replay), so a table answer carries the same
receipts as the `/call/` answer would."""
headers = [(k, v) for k, v in upstream.raw_headers if k.lower().startswith(b"x-treg-")]
status = upstream.status
raw_body, truncated = await read_answer(upstream)
raw_body = _decoded(raw_body, _header(upstream, "content-encoding"))
text = raw_body.decode("utf-8", "replace")
meta = {"call_id": context.call_ref, "cost_micro": context.cost_micro}
if not 200 <= status < 300:
# The provider refused or failed: the same status as /call/ reports, with the start of its
# body, so the add-on can show why. The call's money is already settled (or released).
return status, {"error": "upstream_error", "upstream_status": status,
"body_excerpt": text[:BODY_EXCERPT_BYTES], "_treg": meta}, headers
if truncated:
table = table_domain.raw(text, truncated=True)
else:
try:
body = json.loads(text) if text.strip() else None
except ValueError:
body = None
if body is None:
table = table_domain.raw(text)
else:
tool_id = _tool_id(context, rest)
contract = _contract(tool_id)
if contract is not None:
table = table_domain.to_table(body, contract_output=contract[0], list_field=contract[1])
elif tool_id not in catalog_store.load().by_id and (fields := await _hub_fields(context, rest)) is not None:
table = table_domain.to_table(body, hub_fields=fields)
else:
table = table_domain.to_table(body)
table["_treg"] = {**(table.get("_treg") or {}), **meta}
return status, table, headers
+1
View File
@@ -339,6 +339,7 @@ _DATAPLANE_ROUTE_KEYS: frozenset[RouteKey] = frozenset({
("/catalog/call/{rest:path}",
("DELETE", "GET", "HEAD", "OPTIONS", "PATCH", "POST", "PUT"),
"call_catalog_endpoint"),
("/table/{rest:path}", ("DELETE", "GET", "HEAD", "OPTIONS", "PATCH", "POST", "PUT"), "table_tool"),
# MCP is calling traffic, so its mount and its RFC 9728 resource metadata belong to the
# dataplane. Token issuance (consent, /oauth/*) stays on control; the dataplane only validates.
('/.well-known/oauth-protected-resource/mcp', ('GET',), 'oauth_protected_resource'),
+1
View File
@@ -6,6 +6,7 @@ from __future__ import annotations
_CALL_SURFACES = (
("/catalog/call/", "catalog_call"),
("/call/", "call"),
("/table/", "table"),
)
+14
View File
@@ -465,6 +465,12 @@ class Settings(BaseSettings):
# beside `hub_teams` (owner, 2026-09-26): colleagues try it from their own accounts, with no
# shared team. Either list lets a reader in; both empty means every team.
hub_users: str = ""
# `/table/<tool id>`: the same call as `/call/`, answered as rows and columns (for the Google
# Sheets add-on). Off by default; with it on, `table_teams` / `table_users` limit it the same
# way `hub_teams` / `hub_users` limit the hub, and both empty means every team.
table_enabled: bool = False
table_teams: str = ""
table_users: str = ""
# Additive Claude directory MCP. Default OFF so deploying code cannot publish a new connector
# surface before its production Inspector and custom-connector gates have passed.
@@ -656,6 +662,14 @@ class Settings(BaseSettings):
"""The allow-listed tier-4 providers (comma-separated `TREG_PLATFORM_PROVIDERS`)."""
return frozenset(p.strip().lower() for p in self.platform_providers.split(",") if p.strip())
@property
def table_team_set(self) -> frozenset[str]:
return frozenset(p.strip().lower() for p in self.table_teams.split(",") if p.strip())
@property
def table_user_set(self) -> frozenset[str]:
return frozenset(p.strip().lower() for p in self.table_users.split(",") if p.strip())
@property
def hub_user_set(self) -> frozenset[str]:
"""`hub_users` parsed: lower-cased emails, empty means no restriction by person."""
+277
View File
@@ -0,0 +1,277 @@
"""Turn one treg answer into rows and columns (`POST /table/<tool id>`, for spreadsheets).
Pure and stdlib-only: the parsed JSON answer goes in, with what the catalog knows about the tool,
and a table comes out. No HTTP, no database, no money: the call is already made and settled by the
time this runs (`application/table.py`). Design: tools-gsheet `docs/decisions.md` round 1 q1,
round 2 q9, round 3 B and F.
Four kinds of tool, four sources of columns (`column_source`):
- a routed job with a flat contract output (people.email.verify): the contract's fields, in the
contract's order, then `served_by`; one row, none on a miss;
- a routed job whose contract output has a required list (people.search -> `people`): one row per
item; items are the provider's own objects, so a people-shaped list maps common names to fixed
columns first (`LIST_MAPS`, data, not code), then every other field under its own name;
- a hub tool: its manifest's output fields, one row;
- anything else (`generated`): the lists of objects found in the answer decide the shape: none is
`flat`, one is `list`, two or more are `nested` (`tables` + `summary`).
"""
from __future__ import annotations
import json
from typing import Any
MAX_DEPTH = 3 # a nested object flattens to `a.b.c` columns, no deeper
MAX_WALK = 4 # how deep the search for lists of objects goes
RAW_TEXT_BYTES = 64 * 1024 # the most text a `raw` table carries
# A list holding exactly one object under one of these names is a wrapper, not a table: the search
# for tables goes inside it (seranking `summary[0]`, DataForSEO `tasks[0].result[0]`). Data, so a
# provider's wrapper is one word here, not a code branch.
ENVELOPE_KEYS = frozenset({"summary", "tasks", "result", "results", "data", "response", "response_data"})
# Provider-native rows of a routed list job, mapped to fixed columns first. Each entry: the column,
# then the item paths that may hold it, first match wins. Keyed by the contract's list field.
_PEOPLE_MAP: tuple[tuple[str, tuple[str, ...]], ...] = (
("first_name", ("first_name", "firstname", "firstName")),
("last_name", ("last_name", "lastname", "lastName", "last_name_obfuscated")),
("title", ("title", "job_title", "headline", "lastJobTitle")),
("company", ("company", "company_name", "organization.name", "lastCompanyName")),
("linkedin_url", ("linkedin_url", "linkedinUrl", "profileUrl", "linkedin")),
("location", ("location", "city", "address")),
)
LIST_MAPS: dict[str, tuple[tuple[str, tuple[str, ...]], ...]] = {"people": _PEOPLE_MAP}
# ------------------------------------------------------------------------------------------------
# Cells and records
def cell(value: Any) -> Any:
"""One value as a spreadsheet cell: a scalar as it is; a list of scalars joined with ", "
(decision B); an object, or a list holding objects, as compact JSON."""
if value is None or isinstance(value, (bool, int, float, str)):
return value
if isinstance(value, list) and all(v is None or isinstance(v, (bool, int, float, str)) for v in value):
return ", ".join("" if v is None else str(v) for v in value)
return json.dumps(value, ensure_ascii=False, separators=(",", ":"))
def _hidden(key: Any) -> bool:
return isinstance(key, str) and key.startswith("_")
def flatten(obj: Any, prefix: str = "", depth: int = 0) -> dict[str, Any]:
"""An object as `{column: cell}`: nested objects become `a.b` columns down to MAX_DEPTH, and a
field whose name starts with `_` (treg's own `_treg`, a provider's `_note`) is never a column."""
if not isinstance(obj, dict):
return {prefix or "value": cell(obj)}
out: dict[str, Any] = {}
for key, value in obj.items():
if _hidden(key):
continue
name = f"{prefix}.{key}" if prefix else str(key)
if isinstance(value, dict) and value and depth + 1 < MAX_DEPTH:
out.update(flatten(value, name, depth + 1))
else:
out[name] = cell(value)
return out
def _get_path(obj: Any, path: str) -> tuple[bool, Any]:
cur = obj
for part in path.split("."):
if not isinstance(cur, dict) or part not in cur:
return False, None
cur = cur[part]
return True, cur
def _grid(records: list[dict[str, Any]], fixed: list[str] | None = None) -> tuple[list[str], list[list[Any]]]:
"""Records as columns and rows: the fixed columns first, then every other key in first-seen order."""
columns: list[str] = list(fixed or [])
seen = set(columns)
for rec in records:
for key in rec:
if key not in seen:
seen.add(key)
columns.append(key)
return columns, [[rec.get(c) for c in columns] for rec in records]
# The saved example answers (`src/treg/catalog/examples/`) end a long list with this marker string,
# written by `scripts/catalog_verify.py`; a real answer never has it. Skipped, so an example reads
# like the answer it was cut from.
_TRUNCATION_MARK = " more item(s) truncated"
def _items(value: list) -> list:
return [v for v in value if not (isinstance(v, str) and v.endswith(_TRUNCATION_MARK))]
def _is_record_list(value: Any) -> bool:
if not isinstance(value, list):
return False
items = _items(value)
return bool(items) and all(isinstance(v, dict) for v in items)
# ------------------------------------------------------------------------------------------------
# The four kinds of tool
def from_contract(body: Any, output_fields: list[str], list_field: str | None) -> dict[str, Any]:
"""A routed job's answer, `{output, raw, _treg}`. `list_field` names the contract's required
list, when it has one."""
body = body if isinstance(body, dict) else {}
output = body.get("output") if isinstance(body.get("output"), dict) else {}
meta = body.get("_treg") if isinstance(body.get("_treg"), dict) else {}
served_by = meta.get("served_by")
miss = meta.get("outcome") == "miss"
if list_field:
items = output.get(list_field) if not miss else None
items = _items(items) if isinstance(items, list) else []
mapping = LIST_MAPS.get(list_field)
records = [_mapped(item, mapping) if isinstance(item, dict) else {"value": cell(item)} for item in items]
columns, rows = _grid(records, [c for c, _ in mapping] if mapping else None)
return {"shape": "list", "columns": columns, "rows": rows, "column_source": "contract", "_treg": meta}
columns = list(output_fields) + ["served_by"]
rows = [] if miss else [[cell(output.get(f)) for f in output_fields] + [served_by]]
return {"shape": "flat", "columns": columns, "rows": rows, "column_source": "contract", "_treg": meta}
def _mapped(item: dict[str, Any], mapping) -> dict[str, Any]:
"""One provider row: the mapped columns first (first path that holds a value wins), then every
other field under the provider's own name, minus the fields a mapped column already used."""
rec: dict[str, Any] = {}
used: set[str] = set()
for column, paths in mapping or ():
rec[column] = None
for path in paths:
found, value = _get_path(item, path)
if found and value not in (None, ""):
rec[column] = cell(value)
used.add(path)
break
for key, value in flatten(item).items():
if key not in used and key not in rec:
rec[key] = value
return rec
def from_hub(body: Any, output_fields: list[str]) -> dict[str, Any]:
"""A hub tool's answer, `{run_id, output, usage, trace}`: its manifest's fields, one row."""
body = body if isinstance(body, dict) else {}
output = body.get("output") if isinstance(body.get("output"), dict) else {}
fields = list(output_fields) or list(output)
meta = {"run_id": body.get("run_id"), "recipe": body.get("recipe"), "usage": body.get("usage")}
return {"shape": "flat", "columns": fields, "rows": [[cell(output.get(f)) for f in fields]],
"column_source": "hub", "_treg": meta}
def generated(body: Any) -> dict[str, Any]:
"""Any other answer: the lists of objects in it decide the shape."""
tables = _find_tables(body, "", 0)
if not tables:
base = body
if isinstance(body, list):
# a bare list of scalars, or of mixed values: one column
return {"shape": "list", "columns": ["value"], "rows": [[cell(v)] for v in body],
"column_source": "generated"}
columns, rows = _grid([flatten(base)])
return {"shape": "flat", "columns": columns, "rows": rows, "column_source": "generated"}
if len(tables) == 1:
_, path, items = tables[0]
columns, rows = _grid([flatten(i) for i in items])
return {"shape": "list", "columns": columns, "rows": rows, "column_source": "generated", "path": path}
views = []
for name, path, items in tables:
columns, rows = _grid([flatten(i) for i in items])
views.append({"name": name, "path": path, "row_count": len(rows), "columns": columns, "rows": rows})
summary = _summary(_common_parent(body, [p for _, p, _ in tables]))
return {"shape": "nested", "columns": summary["columns"], "rows": summary["rows"],
"tables": views, "summary": summary, "column_source": "generated"}
def _find_tables(obj: Any, path: str, depth: int) -> list[tuple[str, str, list[dict]]]:
"""Every list of objects in the answer, as (name, path, items). A one-object list under a
wrapper name (`ENVELOPE_KEYS`) is searched inside instead of counted, and a table's own items
are never searched: their inner lists are cells."""
found: list[tuple[str, str, list[dict]]] = []
if depth > MAX_WALK:
return found
if isinstance(obj, list) and _is_record_list(obj) and not path:
return [("rows", "", _items(obj))]
if not isinstance(obj, dict):
return found
for key, value in obj.items():
if _hidden(key):
continue
here = f"{path}.{key}" if path else str(key)
if _is_record_list(value):
items = _items(value)
if len(items) == 1 and key in ENVELOPE_KEYS:
found += _find_tables(items[0], f"{here}[0]", depth + 1)
else:
found.append((str(key), here, items))
elif isinstance(value, dict):
found += _find_tables(value, here, depth + 1)
return found
def _common_parent(body: Any, paths: list[str]) -> Any:
"""The object that holds the tables, when they share one parent (seranking `summary[0]`), so the
summary lists that object's own fields; else the whole answer."""
parents = {p.rsplit(".", 1)[0] if "." in p else "" for p in paths}
if len(parents) != 1:
return body
parent = parents.pop()
cur = body
for part in [p for p in parent.split(".") if p]:
key, _, index = part.partition("[")
cur = cur.get(key) if isinstance(cur, dict) else None
if index:
i = int(index.rstrip("]"))
cur = cur[i] if isinstance(cur, list) and len(cur) > i else None
return cur if isinstance(cur, dict) else body
def _summary(obj: Any) -> dict[str, Any]:
"""An object as two columns, field and value. A list of objects becomes one cell: the first
value of each item, joined with ", " (`top_tlds` -> "com, info, org")."""
rows: list[list[Any]] = []
if isinstance(obj, dict):
for key, value in obj.items():
if _hidden(key):
continue
if _is_record_list(value):
firsts = [next(iter(v.values()), None) for v in _items(value)]
rows.append([str(key), cell([f if not isinstance(f, (dict, list)) else cell(f) for f in firsts])])
elif isinstance(value, dict):
rows += [[k, v] for k, v in flatten(value, str(key)).items()]
else:
rows.append([str(key), cell(value)])
return {"columns": ["field", "value"], "rows": rows}
def raw(text: str, *, truncated: bool = False) -> dict[str, Any]:
"""An answer that is not JSON, or too big to read: its text, cut to RAW_TEXT_BYTES."""
data = text.encode("utf-8", "replace")
cut = len(data) > RAW_TEXT_BYTES
body = data[:RAW_TEXT_BYTES].decode("utf-8", "ignore") if cut else text
return {"shape": "raw", "columns": ["text"], "rows": [[body]], "column_source": "generated",
"truncated": truncated or cut}
def to_table(body: Any, *, contract_output: list[str] | None = None, list_field: str | None = None,
hub_fields: list[str] | None = None) -> dict[str, Any]:
"""The table for one parsed answer. Pass `contract_output` for a routed job (with `list_field`
when its contract has a required list), `hub_fields` for a hub tool, neither for anything else.
Never raises on a strange body: anything it cannot shape becomes a `raw` table."""
try:
if contract_output is not None:
return from_contract(body, contract_output, list_field)
if hub_fields is not None:
return from_hub(body, hub_fields)
return generated(body)
except Exception: # noqa: BLE001 - a table is a view; the answer itself was already delivered
return raw(json.dumps(body, ensure_ascii=False, default=str))
+34 -17
View File
@@ -288,6 +288,37 @@ async def call_tool(
request: Request,
caller: Caller = Depends(require_member),
):
prefix = "/catalog/call/" if getattr(request.state, "catalog_only", False) else "/call/"
return await run_call_surface(rest, request, caller, prefix=prefix, finish=_relay_answer)
async def _relay_answer(request: Request, context, upstream: UpstreamResponse, rest: str) -> Response:
"""`/call/`'s ending: the provider's answer, streamed unchanged, with the optional invitation."""
_attach_async_descriptor(upstream, context, rest)
response = _http_upstream_response(upstream)
try:
# Optional invitation, decided before the body streams (application/call/invite.py).
kind = await invitation(context, response.status_code,
replayed=bool(response.headers.get("X-Treg-Idempotent-Replay")))
if kind is not None:
response.headers["X-Treg-Hint"] = kind
if kind == "review":
response.headers["X-Treg-Review"] = "requested" # read by CLI <= 0.18
_capture_hint(request, context, kind)
except Exception:
pass # A fault here can only lose the header; the answer is already built.
return response
async def run_call_surface(rest: str, request: Request, caller: Caller, *, prefix: str, finish,
headers=None) -> Response:
"""One call through the whole call road: the caller's identity stashed for the refusal fallback,
the context built from the raw request, `execute_call`, and the same bookkeeping on every exit.
Every call surface (`/call/`, `/catalog/call/`, `/table/`) goes through here; only `finish`,
which turns the answer into the HTTP response, differs. `finish(request, context, upstream,
rest)` runs inside the same try, so a fault there is audited and released like any other.
`headers`, when given, replaces the request's raw headers for the upstream call (`/table/` asks
the provider for uncompressed bytes); `/call/` never passes it."""
# Identity for the refusal fallback in `_mark_treg_own_errors`: a raise anywhere below (unknown
# tool, deny rule, daily cap) leaves this handler without an audit row, and the exception handler
# is the one place every such refusal passes through — but it has no Caller of its own.
@@ -307,14 +338,13 @@ async def call_tool(
# valid percent-escapes, so the original bytes travel through to the upstream one-to-one.
raw_path = request.scope.get("raw_path")
if raw_path:
call_prefix = "/catalog/call/" if getattr(request.state, "catalog_only", False) else "/call/"
_, sep, raw_rest = raw_path.decode("ascii", "replace").partition(call_prefix)
_, sep, raw_rest = raw_path.decode("ascii", "replace").partition(prefix)
if sep:
rest = raw_rest
call_input = CallInput(
method=request.method,
raw_rest=rest,
raw_headers=tuple(request.headers.raw),
raw_headers=tuple(request.headers.raw) if headers is None else tuple(headers),
query_items=tuple(request.query_params.multi_items()),
raw_query=request.url.query or "",
body=_HttpRequestBody(request),
@@ -330,20 +360,7 @@ async def call_tool(
request.state.call_context = context
try:
upstream = await execute_call(context, request.app.state.http)
_attach_async_descriptor(upstream, context, rest)
response = _http_upstream_response(upstream)
try:
# Optional invitation, decided before the body streams (application/call/invite.py).
kind = await invitation(context, response.status_code,
replayed=bool(response.headers.get("X-Treg-Idempotent-Replay")))
if kind is not None:
response.headers["X-Treg-Hint"] = kind
if kind == "review":
response.headers["X-Treg-Review"] = "requested" # read by CLI <= 0.18
_capture_hint(request, context, kind)
except Exception:
pass # A fault here can only lose the header; the answer is already built.
return response
return await finish(request, context, upstream, rest)
except CallFailure as exc:
raise _translate_call_failure(exc) from exc
except PoolTimeoutError:
+45
View File
@@ -0,0 +1,45 @@
"""HTTP adapter for `/table/<tool id>`: the same call as `/call/<tool id>`, answered as rows and columns.
Only the ending differs from `/call/`: the call runs through `run_call_surface` (the same gates, key
injection, hold and settle, audit row and `Idempotency-Key`), then `application.table` reads the
answer and converts it. Behind `TREG_TABLE_ENABLED` (+ `TREG_TABLE_TEAMS` / `TREG_TABLE_USERS`).
"""
from __future__ import annotations
from fastapi import APIRouter, Depends, Request
from fastapi.responses import JSONResponse, Response
from ..application import table as table_app
from ..domain.identity.access import Caller, require_member
from .call import run_call_surface
app = APIRouter()
router = app
async def _table_answer(request: Request, context, upstream, rest: str) -> Response:
status, body, headers = await table_app.table_answer(context, upstream, rest)
response = JSONResponse(body, status_code=status)
for name, value in headers:
response.headers[name.decode("latin-1")] = value.decode("latin-1")
return response
@app.api_route(
"/table/{rest:path}",
methods=["GET", "POST", "PUT", "PATCH", "DELETE", "HEAD", "OPTIONS"],
include_in_schema=False,
)
async def table_tool(
rest: str,
request: Request,
caller: Caller = Depends(require_member),
):
"""One call, answered as `{shape, columns, rows, ...}` (docs/context/architecture/table.md)."""
if not table_app.enabled_for(caller.org.slug, caller.email):
# Answered here, not raised: a raised 404 on a call surface is stamped and audited as a
# refusal, and with the flag off this route must look like it does not exist.
return JSONResponse({"detail": "Not Found"}, status_code=404)
return await run_call_surface(rest, request, caller, prefix="/table/", finish=_table_answer,
headers=table_app.plain_headers(tuple(request.headers.raw)))
+14
View File
@@ -2430,6 +2430,20 @@
"name": "call_catalog_endpoint",
"path": "/catalog/call/{rest:path}"
},
{
"kind": "APIRoute",
"methods": [
"DELETE",
"GET",
"HEAD",
"OPTIONS",
"PATCH",
"POST",
"PUT"
],
"name": "table_tool",
"path": "/table/{rest:path}"
},
{
"kind": "APIRoute",
"methods": [
+2
View File
@@ -16,6 +16,7 @@ from treg.bootstrap import ROLE_STARTUP_CHECKS, create_app
_SNAPSHOT = Path(__file__).parent / "snapshots" / "routes.json"
_CALL_ROUTE = "DELETE,GET,HEAD,OPTIONS,PATCH,POST,PUT /call/{rest:path}"
_CATALOG_CALL_ROUTE = "DELETE,GET,HEAD,OPTIONS,PATCH,POST,PUT /catalog/call/{rest:path}"
_TABLE_ROUTE = "DELETE,GET,HEAD,OPTIONS,PATCH,POST,PUT /table/{rest:path}"
def _all_routes() -> list[str]:
@@ -34,6 +35,7 @@ _ALL_ROUTES = _all_routes()
_DATAPLANE_ONLY = (
_CALL_ROUTE,
_CATALOG_CALL_ROUTE,
_TABLE_ROUTE,
"GET,HEAD /.well-known/oauth-protected-resource/mcp",
"GET,HEAD /.well-known/oauth-protected-resource/mcp/v2",
"MOUNT /mcp/v2",
+2
View File
@@ -38,6 +38,8 @@ EXPECTED_MAKERS: dict[str, set[str]] = {
"application/provider_resources.py": {API},
# Interactive paid runs: short transactions between legs, never across upstream waits.
"application/arena.py": {API},
# `/table/`: one short read of a hub tool's manifest, after the call's answer is fully read.
"application/table.py": {API},
# Snapshot read on the request path; the collector runs in the `treg-worker` process (see
# `worker.py` below), where the API pool is the only one in use.
"application/arena_insights.py": {API},
+208
View File
@@ -0,0 +1,208 @@
"""`/table/<tool id>`: the same call as `/call/`, answered as rows and columns.
Two halves: the pure converter (`domain.table`) on saved real answers, and the route end to end
(the flag, the four shapes, an upstream error, the money, and an `Idempotency-Key` replay).
Design: tools-gsheet `docs/decisions.md`; the mechanism: docs/context/architecture/table.md.
"""
from __future__ import annotations
import json
from pathlib import Path
import pytest
from httpx import AsyncClient
from tests.test_marketplace_call import _fake_relay, platform_on # noqa: F401 - tier 4 on
from treg.application import table as table_app
from treg.application.call import service as call_service
from treg.config import get_settings
from treg.domain.table import cell, to_table
EXAMPLES = Path(__file__).resolve().parents[1] / "src" / "treg" / "catalog" / "examples"
EP = "tikhub.tiktok.video.comments" # a real catalog GET endpoint with a platform price
def _example(name: str):
return json.loads((EXAMPLES / f"{name}.json").read_text())
# ------------------------------------------------------------------------------------------------
# The converter, on saved real answers
def test_a_flat_answer_is_one_row_of_its_fields():
t = to_table(_example("bounceban.people.email.verify"))
assert t["shape"] == "flat" and t["column_source"] == "generated" and len(t["rows"]) == 1
row = dict(zip(t["columns"], t["rows"][0]))
assert row["email"] == "dev@bounceban.com" and row["result"] == "deliverable" and row["score"] == 99
def test_one_list_is_one_row_per_item_with_nested_fields_as_dotted_columns():
t = to_table(_example("apollo.people.search"))
assert t["shape"] == "list" and t["path"] == "people" and len(t["rows"]) >= 2
assert "organization.name" in t["columns"] and "first_name" in t["columns"]
assert "_note" not in t["columns"] and "_scrubbed" not in t["columns"] # `_` fields are never columns
def test_several_lists_are_nested_tables_and_a_summary_inside_a_one_item_wrapper():
t = to_table(_example("seranking.web.backlinks.summary"))
assert t["shape"] == "nested"
names = {tb["name"]: tb for tb in t["tables"]}
assert set(names) >= {"top_pages_by_backlinks", "top_tlds", "top_countries"}
assert names["top_tlds"]["path"] == "summary[0].top_tlds" and names["top_tlds"]["columns"] == ["tld", "count"]
summary = dict(t["summary"]["rows"])
assert summary["backlinks"] == 259613 and summary["top_tlds"].startswith("com, info")
assert t["columns"] == ["field", "value"] and t["rows"] == t["summary"]["rows"] # a simple client still works
def test_a_routed_flat_job_uses_the_contract_order_then_served_by():
body = {"output": {"score": 99, "valid": True, "status": "valid"}, "raw": {},
"_treg": {"served_by": "trykitt.people.email.verify", "outcome": "hit"}}
t = to_table(body, contract_output=["valid", "status", "score"])
assert t["columns"] == ["valid", "status", "score", "served_by"]
assert t["rows"] == [[True, "valid", 99, "trykitt.people.email.verify"]]
assert t["column_source"] == "contract" and t["_treg"]["served_by"] == "trykitt.people.email.verify"
def test_a_routed_miss_has_no_rows():
t = to_table({"output": {"valid": None}, "_treg": {"outcome": "miss", "served_by": None}},
contract_output=["valid", "status", "score"])
assert t["rows"] == [] and t["_treg"]["outcome"] == "miss"
def test_a_routed_people_list_maps_common_names_first():
people = _example("apollo.people.search")["people"]
t = to_table({"output": {"people": people}, "_treg": {"served_by": "apollo.people.search"}},
contract_output=["people", "total"], list_field="people")
assert t["shape"] == "list"
assert t["columns"][:6] == ["first_name", "last_name", "title", "company", "linkedin_url", "location"]
row = dict(zip(t["columns"], t["rows"][0]))
assert row["first_name"] == people[0]["first_name"]
assert row["last_name"] == people[0]["last_name_obfuscated"] # a mapped provider name
assert row["company"] == people[0]["organization"]["name"] # a mapped dotted path
assert "last_name_obfuscated" not in t["columns"] # used once, not twice
def test_a_hub_tool_is_its_manifest_fields_in_order():
t = to_table({"run_id": "r1", "output": {"count": 2, "rows": [1, 2]}, "usage": {"cost_micro": 5}},
hub_fields=["rows", "count"])
assert t["columns"] == ["rows", "count"] and t["rows"] == [["1, 2", 2]] and t["column_source"] == "hub"
def test_cells():
assert cell(["a", "b", 3]) == "a, b, 3" # decision B: one cell
assert cell({"a": 1}) == '{"a":1}' and cell([{"a": 1}]) == '[{"a":1}]'
assert cell(None) is None and cell(1.5) == 1.5
def test_a_strange_body_never_raises():
assert to_table(["x", 1, None])["shape"] == "list"
assert to_table("text")["shape"] == "flat"
def test_the_contract_of_a_routed_job_is_found_in_the_catalog():
fields, list_field = table_app._contract("treg.people.email.verify")
assert fields[:3] == ["valid", "status", "score"] and list_field is None
fields, list_field = table_app._contract("treg.people.search")
assert list_field == "people"
assert table_app._contract(EP) is None # not a routed job
# ------------------------------------------------------------------------------------------------
# The route, end to end
@pytest.fixture
def table_on(monkeypatch):
monkeypatch.setenv("TREG_TABLE_ENABLED", "1")
get_settings.cache_clear()
yield
get_settings.cache_clear()
async def _money(clients: AsyncClient) -> tuple[int, list]:
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
b = (await clients.get(f"/orgs/{org_id}/balance")).json()
return b["balance_micro"], b.get("holds") or []
async def test_the_flag_off_is_a_404_that_leaves_no_trace(clients: AsyncClient, platform_on):
r = await clients.get(f"/table/{EP}?aweme_id=1")
assert r.status_code == 404 and "x-treg-error" not in r.headers
calls = (await clients.get("/calls", params={"limit": 20})).json()
assert not [c for c in calls if (c.get("path") or "").startswith("/table/")]
async def test_a_team_outside_the_lists_gets_the_same_404(clients: AsyncClient, platform_on, table_on, monkeypatch):
monkeypatch.setenv("TREG_TABLE_USERS", "someone-else@example.com")
get_settings.cache_clear()
assert (await clients.get(f"/table/{EP}?aweme_id=1")).status_code == 404
async def test_each_shape_end_to_end_and_the_same_charge_as_call(clients: AsyncClient, platform_on, table_on, monkeypatch):
cases = {
"flat": (b'{"id": 7, "author": {"name": "a"}}', lambda t: t["columns"] == ["id", "author.name"]),
"list": (b'{"data": [{"a": 1}, {"a": 2}], "total": 2}', lambda t: t["rows"] == [[1], [2]]),
"nested": (b'{"x": [{"a": 1}, {"a": 2}], "y": [{"b": 1}, {"b": 2}], "n": 2}',
lambda t: [tb["name"] for tb in t["tables"]] == ["x", "y"]),
"raw": (b"not json at all", lambda t: t["rows"] == [["not json at all"]]),
}
before, _ = await _money(clients)
r = await clients.get(f"/call/{EP}?aweme_id=0")
assert r.status_code == 200
after, _ = await _money(clients)
one_call = before - after
assert one_call > 0
for i, (shape, (body, check)) in enumerate(cases.items(), start=1):
monkeypatch.setattr(call_service, "relay", _fake_relay(200, body))
before, _ = await _money(clients)
r = await clients.get(f"/table/{EP}?aweme_id={i}")
assert r.status_code == 200, r.text
t = r.json()
assert t["shape"] == shape and check(t), (shape, t)
assert r.headers["x-treg-call-id"] == t["_treg"]["call_id"]
after, holds = await _money(clients)
assert before - after == one_call, shape # charged exactly as /call/ charges
assert int(r.headers["x-treg-cost-micro"]) == one_call
assert holds == [] # no hold left open
async def test_an_upstream_error_keeps_its_status_and_releases_the_hold(clients: AsyncClient, platform_on, table_on, monkeypatch):
monkeypatch.setattr(call_service, "relay", _fake_relay(502, b'{"message": "vendor down"}'))
before, _ = await _money(clients)
r = await clients.get(f"/table/{EP}?aweme_id=err")
assert r.status_code == 502
body = r.json()
assert body["error"] == "upstream_error" and body["upstream_status"] == 502
assert "vendor down" in body["body_excerpt"]
after, holds = await _money(clients)
assert after == before and holds == [] # a vendor failure costs nothing
async def test_an_idempotency_key_replays_the_same_table_for_free(clients: AsyncClient, platform_on, table_on, monkeypatch):
monkeypatch.setattr(call_service, "relay", _fake_relay(200, b'{"data": [{"a": 1}, {"a": 2}]}'))
h = {"Idempotency-Key": "sheet-1-row-7"}
first = await clients.get(f"/table/{EP}?aweme_id=idem", headers=h)
assert first.status_code == 200, first.text
before, _ = await _money(clients)
again = await clients.get(f"/table/{EP}?aweme_id=idem", headers=h)
after, _ = await _money(clients)
assert again.status_code == 200 and after == before
assert again.headers.get("x-treg-idempotent-replay")
assert again.json()["rows"] == first.json()["rows"] and again.json()["columns"] == first.json()["columns"]
async def test_a_refusal_is_marked_as_treg_s_own_error(clients: AsyncClient, table_on):
r = await clients.get("/table/no-such-tool/x")
assert r.status_code == 404 and r.headers.get("x-treg-error") == "1"
async def test_a_hub_tool_is_its_manifest_fields(clients: AsyncClient, platform_on, table_on, monkeypatch):
from tests.test_hub import PER_ITEM, _publish_script_priced
monkeypatch.setenv("TREG_HUB_ENABLED", "1")
get_settings.cache_clear()
tool_id = await _publish_script_priced(clients, monkeypatch, 0.5, PER_ITEM, n=3)
r = await clients.post(f"/table/{tool_id}", json={})
assert r.status_code == 200, r.text
t = r.json()
assert t["column_source"] == "hub" and t["columns"] == ["rows", "count"]
assert t["rows"] == [["1, 1, 1", 3]]