mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
fix(call): close abandoned idempotency claims with a stored 410; cut primary DB load (#748)
* fix(call): close abandoned idempotency claims with a stored 410 instead of 409 for 24h
A served call whose done-marking failed (a DB pool timeout) left its claim pending, and every
retry of that key answered 409 idempotency_in_progress until the 24h window expired.
- A claim is a lease owned by its call: it carries call_ref from the claim, the owner renews it
every 5 minutes while running, and store, release and renewal are fenced on call_ref.
- A lease older than 15 minutes with no open hold and no pending async task under call_ref (and
its call_ref: children) is closed by compare-and-swap on owner and lease timestamp with a
stored terminal 410: idempotency_response_lost with the charge when the owner's ledger shows
one (GET /calls/{id}/result may still have the answer), else idempotency_outcome_unknown.
The key is never run again: a lapsed lease does not prove the owner stopped, and a second run
under the same key is the double charge idempotency exists to prevent. Legacy rows without a
call_ref keep answering 409 until they expire.
- Storing, releasing and renewing a claim retry once on a pool timeout, on a fresh session.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
* perf(db): index the idempotency sweep and slow the arena refresh
Both run on the 1 vCPU primary whose CPU pinned during the pool-saturation windows that
stranded idempotency claims and settlements:
- _claim_idempotent sweeps expired labels on every keyed call, but idempotentcall had only
membership_id, so the DELETE read every label the caller holds (tens of thousands for a
batch caller). (membership_id, expires_at) makes it a range over the expired rows (0052).
- The arena collector re-aggregated the whole 30-day window every 120 s: three multi-second
aggregates plus a retention DELETE over a multi-GB table, about a quarter of DB time.
A 30-day leaderboard refreshes every 30 minutes now.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
ebe652058b
commit
ca7770a3c7
@@ -211,6 +211,7 @@ Regenerate via `scripts/build-map.py`.
|
||||
| `src/treg/alembic/versions/0050_hub_listing.py` | architecture/hub.md |
|
||||
| `src/treg/alembic/versions/0051_context_dev_tool_host.py` | architecture/auth-secrets.md |
|
||||
| `src/treg/alembic/versions/0052_async_task_hit.py` | architecture/data-model.md |
|
||||
| `src/treg/alembic/versions/0053_idempotentcall_membership_expires_index.py` | architecture/data-model.md |
|
||||
| `src/treg/analytics.py` | architecture/data-model.md |
|
||||
| `src/treg/api.py` | architecture/archive.md, architecture/money.md, architecture/multi-tenancy.md, architecture/proxy-model.md, architecture/super-admin.md, interface/api.md, interface/dashboard.md, interface/landing-sandbox.md, interface/seo.md |
|
||||
| `src/treg/application/__init__.py` | architecture/import-boundaries.md |
|
||||
@@ -646,7 +647,7 @@ Regenerate via `scripts/build-map.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`, `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`, `0052_async_task_hit.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/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`, `0053_idempotentcall_membership_expires_index.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`, `0052_async_task_hit.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` |
|
||||
| `architecture/hub.md` | `__init__.py`, `manifest.py`, `refs.py`, `graph.py`, `__init__.py`, `runner.py`, `sandbox.py`, `limits.py`, `health.py`, `hub_sandbox.py`, `hub.py`, `catalog.py`, `web.py`, `hub_gate.py`, `service.py`, `__init__.py`, `mcp.py`, `hub.js`, `HubPage.vue`, `HubRunPage.vue`, `cli.py`, `worker.py`, `models.py`, `index.html`, `skill.md`, `llms.txt`, `0044_hub_tools.py`, `0045_hub_runs.py`, `0046_hubtool_check_result.py`, `0047_hubrun_output.py`, `0048_hubtool_data.py`, `0049_hubtool_listed_public_log.py`, `0050_hub_listing.py`, `run.js`, `run.js`, `run.js`, `test_hub.py`, `test_hub_sandbox.py`, `test_hub_run.py` |
|
||||
| `architecture/import-boundaries.md` | `pyproject.toml`, `ci.yml`, `__init__.py`, `__init__.py`, `access.py`, `authorize.py`, `idempotency.py`, `overflow.py`, `route.py`, `__init__.py`, `intake.py`, `resolve.py`, `reserve.py`, `settle.py`, `evidence.py`, `service.py`, `types.py`, `client_identity.py`, `__init__.py`, `__init__.py`, `access.py`, `budgets.py`, `publicdemo.py`, `teams.py`, `usage.py`, `__init__.py`, `__init__.py`, `authorization.py`, `oauth_flow.py`, `refresh.py`, `__init__.py`, `__init__.py`, `__init__.py`, `__init__.py`, `__init__.py`, `injectors.py`, `relay.py`, `__init__.py`, `limiter.py`, `test_call_architecture.py`, `test_import_lightness.py` |
|
||||
|
||||
@@ -28,6 +28,7 @@ sources:
|
||||
|
||||
- src/treg/alembic/versions/0011_callrecord_archive_link.py
|
||||
- src/treg/alembic/versions/0015_idempotentcall_membership_cascade.py
|
||||
- src/treg/alembic/versions/0053_idempotentcall_membership_expires_index.py
|
||||
- src/treg/alembic/versions/0034_managed_api_keys.py
|
||||
- src/treg/alembic/versions/0035_default_key_generation.py
|
||||
- src/treg/alembic/versions/0036_activity_key_indexes.py
|
||||
@@ -321,6 +322,16 @@ uses this metadata, never the encrypted token's shape.
|
||||
the membership is revoked there is no valid caller that can replay it. `delete_membership` removes
|
||||
it explicitly and the `membership_id` foreign key uses `ON DELETE CASCADE` as the schema backstop
|
||||
(Alembic `0015`), so a cached paid response can never turn token revocation into a 500.
|
||||
A `pending` row is a lease owned by one call: it carries that call's `call_ref` from the claim,
|
||||
the owner renews `created_at` every `IDEMPOTENCY_LEASE_RENEW_S` while it runs, and every store,
|
||||
release and renewal is fenced on `call_ref`. A lease older than `IDEMPOTENCY_STALE_PENDING_S`
|
||||
with no open hold or pending async task under `call_ref` (and its `call_ref:` children) is closed
|
||||
by compare-and-swap (owner and `created_at`) with a stored terminal 410:
|
||||
`idempotency_response_lost` with the charge when the owner's ledger shows one, else
|
||||
`idempotency_outcome_unknown`. The key is never run again: a lapsed lease does not prove its owner
|
||||
stopped. Rows without a `call_ref` (written before this) keep answering 409 until they expire.
|
||||
The per-call expired-label sweep reads `(membership_id, expires_at)` (Alembic `0053`), so its
|
||||
cost is the expired rows, not every label the caller holds.
|
||||
- **`ToolRequest`** - a "the catalog doesn't have X" report (`POST /tool-requests`, open + per-IP
|
||||
rate-limited): `capability` (the headline, ≤200 chars), `query` (the search that came up empty -
|
||||
auto-filled by agents, the dedup/priority signal), `note`, `contact`, `source` (`web` | `cli` |
|
||||
|
||||
@@ -256,7 +256,9 @@ every two minutes while visible, preserves the last successful values after a re
|
||||
shows the last update time. Prices still come from the catalog and team quote.
|
||||
|
||||
`application.arena_insights.drain` is the collector, run by the `treg-worker arena insights` cron
|
||||
(every two minutes; `--max-seconds` bounds a pass and the next run resumes from the cursor). It no
|
||||
(every two minutes; `--max-seconds` bounds a pass and the next run resumes from the cursor). Once
|
||||
caught up it re-aggregates the 30-day window at most every `REFRESH_SECONDS` (30 minutes): each
|
||||
aggregate scans the whole window on the primary, so the cron cadence is not the refresh cadence. It no
|
||||
longer runs inside the web processes: as a lifespan coroutine, every web process (and every extra
|
||||
instance during a deploy) contended for the cursor row and each walked `callrecord` on the
|
||||
database the money path depends on. In the worker process it uses the API pool, the only one open
|
||||
|
||||
@@ -281,6 +281,9 @@ again, and charges nothing. The result says `replayed: true`.
|
||||
|
||||
Only for a genuine retry. Asking the same question again to see what changed is NEW work: use a new
|
||||
key or none, or you will get the old answer back. Reusing one key for a different request is refused.
|
||||
A 409 means the original call is still running: retry shortly. A 410 `idempotency_response_lost` means
|
||||
it was charged but its answer was not kept: try `GET /calls/{call_id}/result`, or use a new key.
|
||||
A 410 `idempotency_outcome_unknown` means its outcome was not recorded: `GET /calls/{call_id}` shows the cost; use a new key.
|
||||
|
||||
Most retries need none of this — a failed call was never billed.
|
||||
|
||||
|
||||
@@ -284,6 +284,9 @@ again, and charges nothing. The result says `replayed: true`.
|
||||
|
||||
Only for a genuine retry. Asking the same question again to see what changed is NEW work: use a new
|
||||
key or none, or you will get the old answer back. Reusing one key for a different request is refused.
|
||||
A 409 means the original call is still running: retry shortly. A 410 `idempotency_response_lost` means
|
||||
it was charged but its answer was not kept: try `GET /calls/{call_id}/result`, or use a new key.
|
||||
A 410 `idempotency_outcome_unknown` means its outcome was not recorded: `GET /calls/{call_id}` shows the cost; use a new key.
|
||||
|
||||
Most retries need none of this — a failed call was never billed.
|
||||
|
||||
|
||||
@@ -265,6 +265,9 @@ again, and charges nothing. The result says `replayed: true`.
|
||||
|
||||
Only for a genuine retry. Asking the same question again to see what changed is NEW work: use a new
|
||||
key or none, or you will get the old answer back. Reusing one key for a different request is refused.
|
||||
A 409 means the original call is still running: retry shortly. A 410 `idempotency_response_lost` means
|
||||
it was charged but its answer was not kept: try `GET /calls/{call_id}/result`, or use a new key.
|
||||
A 410 `idempotency_outcome_unknown` means its outcome was not recorded: `GET /calls/{call_id}` shows the cost; use a new key.
|
||||
|
||||
Most retries need none of this — a failed call was never billed.
|
||||
|
||||
|
||||
@@ -279,6 +279,9 @@ again, and charges nothing. The result says `replayed: true`.
|
||||
|
||||
Only for a genuine retry. Asking the same question again to see what changed is NEW work: use a new
|
||||
key or none, or you will get the old answer back. Reusing one key for a different request is refused.
|
||||
A 409 means the original call is still running: retry shortly. A 410 `idempotency_response_lost` means
|
||||
it was charged but its answer was not kept: try `GET /calls/{call_id}/result`, or use a new key.
|
||||
A 410 `idempotency_outcome_unknown` means its outcome was not recorded: `GET /calls/{call_id}` shows the cost; use a new key.
|
||||
|
||||
Most retries need none of this — a failed call was never billed.
|
||||
|
||||
|
||||
@@ -263,6 +263,9 @@ again, and charges nothing. The result says `replayed: true`.
|
||||
|
||||
Only for a genuine retry. Asking the same question again to see what changed is NEW work: use a new
|
||||
key or none, or you will get the old answer back. Reusing one key for a different request is refused.
|
||||
A 409 means the original call is still running: retry shortly. A 410 `idempotency_response_lost` means
|
||||
it was charged but its answer was not kept: try `GET /calls/{call_id}/result`, or use a new key.
|
||||
A 410 `idempotency_outcome_unknown` means its outcome was not recorded: `GET /calls/{call_id}` shows the cost; use a new key.
|
||||
|
||||
Most retries need none of this — a failed call was never billed.
|
||||
|
||||
|
||||
@@ -0,0 +1,75 @@
|
||||
"""composite (membership_id, expires_at) on idempotentcall - the per-call expired-label sweep stops
|
||||
reading every label the caller ever stored
|
||||
|
||||
Revision ID: 0053
|
||||
Revises: 0052
|
||||
Create Date: 2026-09-29
|
||||
|
||||
`_claim_idempotent` runs `DELETE ... WHERE membership_id = ? AND expires_at < now` on EVERY call
|
||||
that sends an Idempotency-Key, inside the intake transaction on an api-pool connection. The table
|
||||
carried `membership_id` alone (plus the unique `(membership_id, key)`), so the delete walked every
|
||||
row the caller holds - a stored response body per paid call - to find the handful past their
|
||||
window. A batch caller holds tens of thousands of live labels, and the sweep read all of them per
|
||||
call: 88 ms mean over 3M calls in `pg_stat_statements`, on the same 1 vCPU database whose CPU pinned
|
||||
during the saturation windows that stranded idempotency claims and settlements.
|
||||
|
||||
`(membership_id, expires_at)` turns the sweep into a range over only the expired rows. The existing
|
||||
single-column index stays: dropping it is a separate, non-additive change.
|
||||
|
||||
Built with the 0020/0021 discipline (CONCURRENTLY in an autocommit block, INVALID debris dropped
|
||||
first). The expand-safety linter counts the autocommit escape as non-additive, so this revision
|
||||
declares a rollback floor pro forma: the operation is one additive index.
|
||||
"""
|
||||
from collections.abc import Sequence
|
||||
|
||||
import sqlalchemy as sa
|
||||
from alembic import op
|
||||
|
||||
revision: str = "0053"
|
||||
down_revision: str | Sequence[str] | None = "0052"
|
||||
branch_labels: str | Sequence[str] | None = None
|
||||
depends_on: str | Sequence[str] | None = None
|
||||
contract = True # pro forma - see the rollback floor note; the operation is one additive index
|
||||
|
||||
_TABLE = "idempotentcall"
|
||||
_INDEXES = (
|
||||
("ix_idempotentcall_membership_id_expires_at", ["membership_id", "expires_at"]),
|
||||
)
|
||||
|
||||
# Same values and reasoning as 0021; `env.py`'s values are restored before the block ends.
|
||||
_LOCK_TIMEOUT = "180s"
|
||||
_STATEMENT_TIMEOUT = "600s"
|
||||
_ENV_LOCK_TIMEOUT = "5s"
|
||||
_ENV_STATEMENT_TIMEOUT = "120s"
|
||||
|
||||
_VALIDITY = sa.text(
|
||||
"SELECT i.indisvalid FROM pg_class c JOIN pg_index i ON i.indexrelid = c.oid "
|
||||
"WHERE c.relname = :name")
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
if op.get_bind().dialect.name != "postgresql":
|
||||
for name, columns in _INDEXES: # SQLite: no concurrent mode, no traffic to block
|
||||
op.create_index(name, _TABLE, columns)
|
||||
return
|
||||
# CONCURRENTLY cannot run inside a transaction; alembic opens one by default.
|
||||
with op.get_context().autocommit_block():
|
||||
bind = op.get_bind()
|
||||
bind.execute(sa.text(f"SET lock_timeout = '{_LOCK_TIMEOUT}'"))
|
||||
bind.execute(sa.text(f"SET statement_timeout = '{_STATEMENT_TIMEOUT}'"))
|
||||
try:
|
||||
for name, columns in _INDEXES:
|
||||
valid = bind.execute(_VALIDITY, {"name": name}).scalar()
|
||||
if valid is True:
|
||||
continue
|
||||
if valid is False: # debris from a killed build - unusable, and never repaired
|
||||
op.drop_index(name, table_name=_TABLE, postgresql_concurrently=True)
|
||||
op.create_index(name, _TABLE, columns, postgresql_concurrently=True)
|
||||
finally:
|
||||
bind.execute(sa.text(f"SET lock_timeout = '{_ENV_LOCK_TIMEOUT}'"))
|
||||
bind.execute(sa.text(f"SET statement_timeout = '{_ENV_STATEMENT_TIMEOUT}'"))
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
for name, _ in _INDEXES:
|
||||
op.drop_index(name, table_name=_TABLE)
|
||||
@@ -29,7 +29,9 @@ from ..models import ArenaInsightState, ArenaObservation, ArchiveKey, ArchiveSna
|
||||
from ..timeutil import utcnow_naive as now
|
||||
|
||||
WINDOW_DAYS = 30
|
||||
REFRESH_SECONDS = 120
|
||||
# Each refresh re-aggregates the whole 30-day window on the primary (tens of seconds on a large
|
||||
# table, evicting the money path's cached pages). A 30-day leaderboard does not need minutes.
|
||||
REFRESH_SECONDS = 1800
|
||||
BATCH_SIZE = 100
|
||||
MAX_BODY_BYTES = 2_000_000
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
@@ -2,20 +2,21 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from sqlalchemy import delete
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy import delete, func, or_, update
|
||||
from sqlalchemy.exc import IntegrityError, TimeoutError as PoolTimeoutError
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlmodel import select
|
||||
|
||||
from ...infra.db import session_maker
|
||||
from ...domain.identity.access import Caller
|
||||
from ...models import IdempotentCall
|
||||
from ...models import AsyncTaskRecord, Hold, IdempotentCall, LedgerEntry
|
||||
from .types import IdempotencyFailed, IdempotentReplay
|
||||
|
||||
if TYPE_CHECKING:
|
||||
@@ -24,6 +25,15 @@ if TYPE_CHECKING:
|
||||
|
||||
IDEMPOTENCY_WINDOW_S = 24 * 3600 # retries happen in seconds; a day is generous and easy to reason about
|
||||
IDEMPOTENCY_HEADER = "idempotency-key"
|
||||
# A pending claim is a LEASE owned by one call (`call_ref`). The owner renews it every
|
||||
# IDEMPOTENCY_LEASE_RENEW_S while it runs (`_hold_claim_lease`), so `created_at` older than
|
||||
# IDEMPOTENCY_STALE_PENDING_S means the owner stopped: its process died, or the bookkeeping that
|
||||
# marks the claim done or releases it failed (a pool timeout). Such a claim is closed with a stored
|
||||
# terminal answer rather than answering 409 for the whole window; see `_resolve_stale_claim`.
|
||||
IDEMPOTENCY_LEASE_RENEW_S = 5 * 60
|
||||
IDEMPOTENCY_STALE_PENDING_S = 3 * IDEMPOTENCY_LEASE_RENEW_S
|
||||
_POOL_RETRY_PAUSE_S = 0.5 # the settle retry's pause (settle.py); one retry, fresh session
|
||||
|
||||
_IDEM_MAX_KEY = 200
|
||||
|
||||
|
||||
@@ -119,6 +129,11 @@ async def _replay_idempotent(key: str, fingerprint: str, caller: Caller,
|
||||
"idempotency_mismatch", status_code=422,
|
||||
detail=(f"Idempotency-Key {_idem_display(key)!r} was already used for a different request. Use a new key, or "
|
||||
f"repeat the original request exactly."))
|
||||
if (row.status != "done" and row.call_ref
|
||||
and row.created_at < _utcnow() - timedelta(seconds=IDEMPOTENCY_STALE_PENDING_S)):
|
||||
resolved = await _resolve_stale_claim(row, caller, db)
|
||||
if resolved is not None:
|
||||
return resolved
|
||||
if row.status != "done" or row.response_status is None:
|
||||
# Still in flight. The first call is talking to the provider right now; telling the caller to
|
||||
# retry is honest and cheap, and it is what stops the second one duplicating the spend.
|
||||
@@ -137,7 +152,134 @@ async def _replay_idempotent(key: str, fingerprint: str, caller: Caller,
|
||||
)
|
||||
|
||||
|
||||
async def _release_idempotent_claim(claim: tuple[int, str] | None) -> None:
|
||||
def _utcnow() -> datetime:
|
||||
return datetime.now(timezone.utc).replace(tzinfo=None)
|
||||
|
||||
|
||||
def _owned(membership_id: int, key: str, call_ref: str):
|
||||
"""The WHERE clause every owner write carries: a claim taken over by another call is not ours."""
|
||||
return (IdempotentCall.membership_id == membership_id, IdempotentCall.key == key,
|
||||
IdempotentCall.call_ref == call_ref)
|
||||
|
||||
|
||||
async def _money_of(db: AsyncSession, org_id: int, call_ref: str, since: datetime,
|
||||
until: datetime) -> dict:
|
||||
"""What one call and every child it spawned did with money, from the hold and ledger tables.
|
||||
|
||||
A call's money lives under ids prefixed by its call_ref: its own hold, `:overflow`, routed
|
||||
children `:r{n}` and their polls, hub steps `:s{n}` and the hub `:price`. uuid hex carries no
|
||||
LIKE wildcard. The LIKE cannot use an index under a non-C collation, so the ledger read is
|
||||
bounded to the owner's lifetime on `(org_id, created_at)`: from the claim to one renewal past
|
||||
its last lease renewal. The only money that can move after that is an async task's worker
|
||||
settling or releasing it, so the tasks the owner started in its lifetime are found the same
|
||||
way and their ledger rows read by exact call id, with no time bound: a still-pending task is
|
||||
money in flight, and a late settle is a charge.
|
||||
"""
|
||||
def mine(col):
|
||||
return or_(col == call_ref, col.like(f"{call_ref}:%"))
|
||||
open_hold = (await db.execute(select(Hold.id).where(
|
||||
Hold.org_id == org_id, mine(Hold.id)).limit(1))).first() is not None
|
||||
rows = list((await db.execute(select(LedgerEntry.kind, LedgerEntry.amount_micro).where(
|
||||
LedgerEntry.org_id == org_id, LedgerEntry.created_at >= since, LedgerEntry.created_at <= until,
|
||||
mine(LedgerEntry.call_id)))).all())
|
||||
tasks = (await db.execute(select(AsyncTaskRecord.call_id, AsyncTaskRecord.status).where(
|
||||
AsyncTaskRecord.org_id == org_id, AsyncTaskRecord.created_at >= since,
|
||||
AsyncTaskRecord.created_at <= until, mine(AsyncTaskRecord.call_id)))).all()
|
||||
if any(status == "pending" for _, status in tasks):
|
||||
open_hold = True # a task's hold is its worker's until it resolves
|
||||
if tasks:
|
||||
rows += (await db.execute(select(LedgerEntry.kind, LedgerEntry.amount_micro).where(
|
||||
LedgerEntry.call_id.in_([task_id for task_id, _ in tasks]),
|
||||
LedgerEntry.created_at > until))).all()
|
||||
# settle entries are negative (money leaving)
|
||||
return {"in_flight": open_hold, "charged": -sum(amount for kind, amount in rows if kind == "settle")}
|
||||
|
||||
|
||||
async def _resolve_stale_claim(row: IdempotentCall, caller: Caller, db: AsyncSession):
|
||||
"""Close an abandoned lease with a stored terminal answer; never run the key again.
|
||||
|
||||
An open hold or a pending async task means the owner or its worker is not done: 409 (None
|
||||
here). Otherwise the answer is a stored 410, so every later retry gets the same reply:
|
||||
`idempotency_response_lost` with the charge when the owner's money shows one (the answer was
|
||||
paid for but not kept), else `idempotency_outcome_unknown`. The claim is never handed to a new
|
||||
call: a lease that stopped renewing does not prove its owner stopped, and a second run under
|
||||
the same key is exactly the double charge this table exists to prevent. A new key calls again.
|
||||
|
||||
The write is a compare-and-swap on the lease as read (owner AND `created_at`): a renewal or
|
||||
another retry in between wins, and this one answers 409 on the fresh state.
|
||||
"""
|
||||
started = row.expires_at - timedelta(seconds=IDEMPOTENCY_WINDOW_S)
|
||||
lifetime_end = row.created_at + timedelta(seconds=IDEMPOTENCY_LEASE_RENEW_S + 60)
|
||||
money = await _money_of(db, caller.org_id, row.call_ref, started, lifetime_end)
|
||||
if money["in_flight"]:
|
||||
return None
|
||||
label, owner, charged = _idem_display(row.key), row.call_ref, max(0, money["charged"])
|
||||
if charged:
|
||||
detail = {"error": "idempotency_response_lost", "call_id": owner, "charged_micro": charged,
|
||||
"message": (f"the call with Idempotency-Key {label!r} completed and was charged, but "
|
||||
f"its response was not retained. GET /calls/{owner}/result may still "
|
||||
"have it; send a new key to call again.")}
|
||||
else:
|
||||
detail = {"error": "idempotency_outcome_unknown", "call_id": owner,
|
||||
"message": (f"the call with Idempotency-Key {label!r} stopped before its outcome was "
|
||||
f"recorded and may have reached the provider. GET /calls/{owner} shows "
|
||||
"what it cost; send a new key to call again.")}
|
||||
body = json.dumps({"detail": detail}, separators=(",", ":")).encode()
|
||||
won = (await db.execute(update(IdempotentCall).where(
|
||||
IdempotentCall.id == row.id, IdempotentCall.status == "pending",
|
||||
IdempotentCall.call_ref == owner, IdempotentCall.created_at == row.created_at,
|
||||
).values(status="done", response_status=410, response_body=body,
|
||||
response_media_type="application/json", charged_micro=charged,
|
||||
# A lost swap must not paint the loaded row as done: the caller re-reads it as a 409.
|
||||
).execution_options(synchronize_session=False))).rowcount == 1
|
||||
await db.commit()
|
||||
if not won:
|
||||
return None
|
||||
logging.getLogger("treg.idempotency").warning(
|
||||
"idempotency claim %s: owner %s abandoned it; stored 410 %s", label, owner, detail["error"])
|
||||
return IdempotentReplay(body=body, status_code=410, media_type="application/json",
|
||||
charged_micro=charged, call_ref=owner)
|
||||
|
||||
|
||||
async def _with_pool_retry(work):
|
||||
"""One retry after a pool timeout, on a fresh session: the settle retry's shape. A saturated
|
||||
pool is the usual reason bookkeeping fails, and a failed write here strands a claim."""
|
||||
try:
|
||||
return await work()
|
||||
except PoolTimeoutError:
|
||||
await asyncio.sleep(_POOL_RETRY_PAUSE_S)
|
||||
return await work()
|
||||
|
||||
|
||||
async def _renew_claim_lease(claim: tuple[int, str, str]) -> None:
|
||||
"""Refresh the owner's lease. Never raises: a missed renewal only shortens the lease."""
|
||||
membership_id, key, call_ref = claim
|
||||
|
||||
async def _renew() -> None:
|
||||
async with session_maker() as db:
|
||||
await db.execute(update(IdempotentCall).where(
|
||||
*_owned(membership_id, key, call_ref), IdempotentCall.status == "pending",
|
||||
).values(created_at=_utcnow()))
|
||||
await db.commit()
|
||||
|
||||
try:
|
||||
await _with_pool_retry(_renew)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logging.getLogger("treg.idempotency").warning(
|
||||
"could not renew idempotency claim %s: %s", _idem_display(key), exc)
|
||||
|
||||
|
||||
async def _hold_claim_lease(state) -> None:
|
||||
"""Renew `state.idem_claim` until the request drops it; the caller cancels this task on exit."""
|
||||
while True:
|
||||
await asyncio.sleep(IDEMPOTENCY_LEASE_RENEW_S)
|
||||
claim = getattr(state, "idem_claim", None)
|
||||
if not claim:
|
||||
return
|
||||
await _renew_claim_lease(claim)
|
||||
|
||||
|
||||
async def _release_idempotent_claim(claim: tuple[int, str, str] | None) -> None:
|
||||
"""Drop a claim this request took and never completed, so the label is usable again at once.
|
||||
|
||||
Does nothing when there is no claim, which is every request that sent no key. Never raises: this
|
||||
@@ -145,27 +287,29 @@ async def _release_idempotent_claim(claim: tuple[int, str] | None) -> None:
|
||||
"""
|
||||
if not claim:
|
||||
return
|
||||
membership_id, key = claim
|
||||
try:
|
||||
membership_id, key, call_ref = claim
|
||||
|
||||
async def _drop() -> None:
|
||||
async with session_maker() as db:
|
||||
row = (await db.execute(select(IdempotentCall).where(
|
||||
IdempotentCall.membership_id == membership_id,
|
||||
IdempotentCall.key == key,
|
||||
IdempotentCall.status == "pending"))).scalar_one_or_none()
|
||||
if row is not None:
|
||||
await db.delete(row)
|
||||
await db.commit()
|
||||
await db.execute(delete(IdempotentCall).where(
|
||||
*_owned(membership_id, key, call_ref), IdempotentCall.status == "pending"))
|
||||
await db.commit()
|
||||
|
||||
try:
|
||||
await _with_pool_retry(_drop)
|
||||
except Exception as exc: # noqa: BLE001 — an error is already on its way out
|
||||
logging.getLogger("treg.idempotency").error(
|
||||
"could not release idempotency claim %s: %s", key, exc, exc_info=True)
|
||||
|
||||
|
||||
async def _claim_idempotent(key: str, fingerprint: str, rest: str, caller: Caller,
|
||||
db: AsyncSession) -> bool:
|
||||
db: AsyncSession, *, call_ref: str) -> bool:
|
||||
"""Take the label for this caller, or report that somebody else already has it.
|
||||
|
||||
The pending row IS the lock. It goes in before the upstream call, so a concurrent retry loses the
|
||||
insert on `(membership_id, key)` and is told to wait rather than duplicating the spend.
|
||||
insert on `(membership_id, key)` and is told to wait rather than duplicating the spend. It
|
||||
carries the owner's `call_ref` from the start: every later write is fenced on it, and a stale
|
||||
claim is resolved from that call's money.
|
||||
"""
|
||||
# Sweep this caller's expired labels first. LAZY and caller-scoped, matching the hold reaper in
|
||||
# domain/money and for the same reasons: a background timer would need a scheduler and a leader
|
||||
@@ -180,7 +324,7 @@ async def _claim_idempotent(key: str, fingerprint: str, rest: str, caller: Calle
|
||||
IdempotentCall.expires_at < datetime.now(timezone.utc).replace(tzinfo=None)))
|
||||
|
||||
row = IdempotentCall(
|
||||
org_id=caller.org_id, membership_id=caller.membership.id, key=key,
|
||||
org_id=caller.org_id, membership_id=caller.membership.id, key=key, call_ref=call_ref,
|
||||
request_fingerprint=fingerprint, endpoint_id=rest[:200], status="pending",
|
||||
expires_at=datetime.now(timezone.utc).replace(tzinfo=None)
|
||||
+ timedelta(seconds=IDEMPOTENCY_WINDOW_S))
|
||||
@@ -195,7 +339,7 @@ async def _claim_idempotent(key: str, fingerprint: str, rest: str, caller: Calle
|
||||
|
||||
async def _store_idempotent(key: str, caller: Caller, *, status_code: int, body: bytes,
|
||||
media_type: str, charged_micro: int, metered: bool,
|
||||
call_ref: str = "", terminal: bool = False) -> None:
|
||||
call_ref: str, terminal: bool = False) -> None:
|
||||
"""Remember a metered success or an explicitly terminal, partially charged routed failure.
|
||||
|
||||
Metered only. A team calling on its OWN key is billed by the provider, not by us, so there is
|
||||
@@ -208,27 +352,25 @@ async def _store_idempotent(key: str, caller: Caller, *, status_code: int, body:
|
||||
wait out the window before they can try again.
|
||||
|
||||
Never raises: the caller already has their answer, and a bookkeeping failure must not turn a
|
||||
served call into a 500. Its own session, because the request's may be mid-rollback.
|
||||
served call into a 500. Its own session, because the request's may be mid-rollback. Fenced on
|
||||
`call_ref`: a claim another call took over is left alone.
|
||||
"""
|
||||
keep = metered and (200 <= status_code < 300 or terminal)
|
||||
try:
|
||||
owned = (*_owned(caller.membership.id, key, call_ref), IdempotentCall.status == "pending")
|
||||
|
||||
async def _write() -> None:
|
||||
async with session_maker() as db:
|
||||
row = (await db.execute(select(IdempotentCall).where(
|
||||
IdempotentCall.membership_id == caller.membership.id,
|
||||
IdempotentCall.key == key))).scalar_one_or_none()
|
||||
if row is None:
|
||||
return
|
||||
if not keep:
|
||||
await db.delete(row)
|
||||
await db.execute(delete(IdempotentCall).where(*owned))
|
||||
else:
|
||||
row.status = "done"
|
||||
row.response_status = status_code
|
||||
row.response_body = body
|
||||
row.response_media_type = media_type or "application/json"
|
||||
row.charged_micro = charged_micro
|
||||
row.call_ref = call_ref
|
||||
db.add(row)
|
||||
await db.execute(update(IdempotentCall).where(*owned).values(
|
||||
status="done", response_status=status_code, response_body=body,
|
||||
response_media_type=media_type or "application/json",
|
||||
charged_micro=charged_micro))
|
||||
await db.commit()
|
||||
|
||||
try:
|
||||
await _with_pool_retry(_write)
|
||||
except Exception as exc: # noqa: BLE001 — loudly, but never into the caller's response
|
||||
logging.getLogger("treg.idempotency").error(
|
||||
"could not record idempotency key %s: %s", key, exc, exc_info=True)
|
||||
|
||||
@@ -126,7 +126,7 @@ class IntakeResult:
|
||||
idempotency_key: str
|
||||
fingerprint: str
|
||||
replay: IdempotentReplay | None
|
||||
claim: tuple[int, str] | None
|
||||
claim: tuple[int, str, str] | None
|
||||
|
||||
|
||||
async def prepare_call_intake(
|
||||
@@ -139,6 +139,7 @@ async def prepare_call_intake(
|
||||
read_body: Callable[[], Awaitable[bytes]],
|
||||
caller: Caller,
|
||||
enforce_tag_budgets: Callable[[Caller, CallMeta, AsyncSession], Awaitable[None]],
|
||||
call_ref: str,
|
||||
) -> IntakeResult:
|
||||
"""Run the pre-resolve gates in their frozen order using bounded transactions."""
|
||||
# Blocked status and the per-tag call count, BEFORE the replay below: a blocked user must neither
|
||||
@@ -167,9 +168,9 @@ async def prepare_call_intake(
|
||||
# miss the lookup above; the unique constraint is what makes the loser wait instead of making
|
||||
# a second upstream call. A check-then-act in Python would leave exactly the window this
|
||||
# feature exists to close — the same reasoning as the conditional UPDATE in ledger.reserve.
|
||||
if not await _claim_idempotent(key, fingerprint, rest, caller, db):
|
||||
if not await _claim_idempotent(key, fingerprint, rest, caller, db, call_ref=call_ref):
|
||||
raise IdempotencyFailed(
|
||||
"idempotency_in_progress", status_code=409,
|
||||
detail=(f"a call with Idempotency-Key {_idem_display(key)!r} "
|
||||
"is already in progress — retry shortly"))
|
||||
return IntakeResult(key, fingerprint, None, (caller.membership.id, key))
|
||||
return IntakeResult(key, fingerprint, None, (caller.membership.id, key, call_ref))
|
||||
|
||||
@@ -41,7 +41,7 @@ from .evidence import (
|
||||
_redact_snippet,
|
||||
_safe_secret_renderings,
|
||||
)
|
||||
from .idempotency import IDEMPOTENCY_HEADER, _store_idempotent
|
||||
from .idempotency import IDEMPOTENCY_HEADER, _hold_claim_lease, _store_idempotent
|
||||
from .intake import META_HEADER, _parse_call_meta, _tag_telemetry, prepare_call_intake
|
||||
from .reserve import _enforce_tag_budgets, _platform_reserve
|
||||
from .resolve import (
|
||||
@@ -104,6 +104,7 @@ class _ApplicationRequest:
|
||||
call_cost_micro=context.cost_micro,
|
||||
)
|
||||
self.db = session_maker()
|
||||
self.lease: asyncio.Task | None = None
|
||||
|
||||
|
||||
def _served_response(served: dict, body: bytes) -> UpstreamResponse:
|
||||
@@ -150,6 +151,8 @@ async def execute_call(context: CallContext, upstream_client: httpx.AsyncClient)
|
||||
try:
|
||||
return await _execute_call(request, upstream_client)
|
||||
finally:
|
||||
if request.lease is not None:
|
||||
request.lease.cancel()
|
||||
context.idempotency = request.state.idem_claim
|
||||
context.audited = request.state.call_audited
|
||||
context.cost_micro = request.state.call_cost_micro
|
||||
@@ -571,6 +574,7 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
|
||||
read_body=request.body,
|
||||
caller=caller,
|
||||
enforce_tag_budgets=_enforce_tag_budgets,
|
||||
call_ref=call_ref,
|
||||
)
|
||||
except CallFailure:
|
||||
raise
|
||||
@@ -586,13 +590,19 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
|
||||
# A replay charges nothing; the first call's charge is echoed separately so a
|
||||
# client summing X-Treg-Cost-Micro never counts one call twice.
|
||||
"X-Treg-Cost-Micro": "0",
|
||||
"X-Treg-Original-Cost-Micro": str(replayed.charged_micro),
|
||||
# Omitted for a stale claim closed as `idempotency_outcome_unknown` (a
|
||||
# stored 410 with no charge found): its original cost is not known.
|
||||
**({} if replayed.status_code == 410 and not replayed.charged_micro
|
||||
else {"X-Treg-Original-Cost-Micro": str(replayed.charged_micro)}),
|
||||
**({"X-Treg-Error": "1"} if replayed.status_code >= 400 else {}),
|
||||
**({"X-Treg-Call-Id": replayed.call_ref} if replayed.call_ref else {})},
|
||||
)
|
||||
# Park it so a failure anywhere below can give the label back. Set AFTER the claim succeeds,
|
||||
# so losing the race above never releases the winner's row.
|
||||
request.state.idem_claim = intake.claim
|
||||
if intake.claim:
|
||||
# Keeps the claim's lease fresh while this call runs; `execute_call` cancels it on exit.
|
||||
request.lease = asyncio.create_task(_hold_claim_lease(request.state))
|
||||
|
||||
drop_params: set[str] = set()
|
||||
streaming_free_result = False
|
||||
@@ -1604,7 +1614,8 @@ async def _execute_call(request: _ApplicationRequest, upstream_client: httpx.Asy
|
||||
# label at once instead of making the caller wait out the window to reuse it.
|
||||
try:
|
||||
await _store_idempotent(idem_key, caller, status_code=response.status, body=b"",
|
||||
media_type="", charged_micro=0, metered=False)
|
||||
media_type="", charged_micro=0, metered=False,
|
||||
call_ref=call_ref)
|
||||
except asyncio.CancelledError:
|
||||
await _finish_cancelled_call(request, mk, call_ref, response)
|
||||
raise
|
||||
|
||||
@@ -1262,7 +1262,7 @@ async def close_deferred(items: list[DeferredSettle], *, charge: bool, why: str
|
||||
|
||||
|
||||
async def _finish_cancelled_call(
|
||||
claim: tuple[int, str] | None,
|
||||
claim: tuple[int, str, str] | None,
|
||||
mk: MarketplaceCall | None,
|
||||
call_ref: str,
|
||||
response: UpstreamResponse | None = None,
|
||||
|
||||
+5
-1
@@ -1239,7 +1239,11 @@ class IdempotentCall(SQLModel, table=True):
|
||||
one that says no.
|
||||
"""
|
||||
|
||||
__table_args__ = (UniqueConstraint("membership_id", "key", name="uq_idem_caller_key"),)
|
||||
__table_args__ = (
|
||||
UniqueConstraint("membership_id", "key", name="uq_idem_caller_key"),
|
||||
# The per-call expired-label sweep in `_claim_idempotent` (Alembic 0053).
|
||||
Index("ix_idempotentcall_membership_id_expires_at", "membership_id", "expires_at"),
|
||||
)
|
||||
|
||||
id: int | None = Field(default=None, primary_key=True)
|
||||
# org_id is kept alongside the caller so the row is still org-scoped for deletion and audit.
|
||||
|
||||
@@ -116,6 +116,13 @@ Keys are yours alone: scoped to the caller, kept 24 hours. Nothing happens witho
|
||||
call that fails frees its key at once. Over MCP the same thing is the `idempotency_key` argument to
|
||||
`call`, and a replayed result carries `replayed: true`.
|
||||
|
||||
A `409` means the original call with that key is still running: retry shortly. A `410`
|
||||
`idempotency_response_lost` means the original call completed and was charged but treg could not keep
|
||||
its answer; the body names the `call_id` and the charge, `GET /calls/{call_id}/result` may still have
|
||||
the answer, and a new key calls again. A `410` `idempotency_outcome_unknown` means the original call
|
||||
stopped before its outcome was recorded and may have reached the provider: `GET /calls/{call_id}`
|
||||
shows what it cost, and a new key calls again.
|
||||
|
||||
### Reselling treg inside your own product
|
||||
|
||||
If you are a platform billing your own users, run one team, one balance, one token, and tag each call
|
||||
|
||||
@@ -264,6 +264,9 @@ again, and charges nothing. The result says `replayed: true`.
|
||||
|
||||
Only for a genuine retry. Asking the same question again to see what changed is NEW work: use a new
|
||||
key or none, or you will get the old answer back. Reusing one key for a different request is refused.
|
||||
A 409 means the original call is still running: retry shortly. A 410 `idempotency_response_lost` means
|
||||
it was charged but its answer was not kept: try `GET /calls/{call_id}/result`, or use a new key.
|
||||
A 410 `idempotency_outcome_unknown` means its outcome was not recorded: `GET /calls/{call_id}` shows the cost; use a new key.
|
||||
|
||||
Most retries need none of this — a failed call was never billed.
|
||||
|
||||
|
||||
@@ -156,7 +156,7 @@ async def test_missing_archive_is_revisited_and_old_window_removed(clients):
|
||||
db.add(key);await db.flush()
|
||||
db.add(ArchiveSnapshot(key_id=key.id,content_hash="late-body",body=b'{"data":{"email":"new@example.test"}}'))
|
||||
row.archive_key_hash=key.key_hash;row.archive_content_hash="late-body";db.add(row)
|
||||
state=(await db.execute(select(ArenaInsightState))).scalar_one();state.updated_at=now()-timedelta(seconds=121)
|
||||
state=(await db.execute(select(ArenaInsightState))).scalar_one();state.updated_at=now()-timedelta(seconds=service.REFRESH_SECONDS+1)
|
||||
db.add(state);await db.commit()
|
||||
await service.collect_batch(session_maker)
|
||||
row=(await service.public_snapshot(session_maker))["rows"][0]
|
||||
|
||||
@@ -227,7 +227,7 @@ async def test_repeated_cancellation_cannot_interrupt_compensation(
|
||||
)
|
||||
original_http = app.state.http
|
||||
original_commit = AsyncSession.commit
|
||||
original_delete = AsyncSession.delete
|
||||
original_execute = AsyncSession.execute
|
||||
original_release = ledger.release_in_transaction
|
||||
cleanup_commit_reached = asyncio.Event()
|
||||
allow_cleanup_commit = asyncio.Event()
|
||||
@@ -239,10 +239,13 @@ async def test_repeated_cancellation_cannot_interrupt_compensation(
|
||||
db.sync_session.info["cancelled_release"] = True
|
||||
return await original_release(db, call_id, **kwargs)
|
||||
|
||||
async def _tag_claim_delete(db: AsyncSession, row) -> None:
|
||||
if isinstance(row, IdempotentCall) and row.key == key:
|
||||
async def _tag_claim_delete(db: AsyncSession, statement, *args, **kwargs):
|
||||
# The release is one DELETE fenced on the claim's owner, not a load-then-delete.
|
||||
if (getattr(statement, "is_delete", False)
|
||||
and statement.table.name == IdempotentCall.__tablename__
|
||||
and "call_ref" in str(statement)): # not the intake's expired-label sweep
|
||||
db.sync_session.info["claim_delete"] = True
|
||||
await original_delete(db, row)
|
||||
return await original_execute(db, statement, *args, **kwargs)
|
||||
|
||||
async def _gate_cleanup_commit(db: AsyncSession) -> None:
|
||||
if db.sync_session.info.get("cancelled_release"):
|
||||
@@ -254,7 +257,7 @@ async def test_repeated_cancellation_cannot_interrupt_compensation(
|
||||
await original_commit(db)
|
||||
|
||||
monkeypatch.setattr(ledger, "release_in_transaction", _tag_cancelled_release)
|
||||
monkeypatch.setattr(AsyncSession, "delete", _tag_claim_delete)
|
||||
monkeypatch.setattr(AsyncSession, "execute", _tag_claim_delete)
|
||||
monkeypatch.setattr(AsyncSession, "commit", _gate_cleanup_commit)
|
||||
app.state.http = tracked
|
||||
headers = {"Idempotency-Key": key}
|
||||
|
||||
@@ -1769,6 +1769,281 @@ async def test_the_sweep_clears_labels_NOBODY_COMES_BACK_FOR(clients: AsyncClien
|
||||
assert gone is None, "an expired row nobody returns for must still be reclaimed"
|
||||
|
||||
|
||||
async def _claim_row(key: str):
|
||||
from sqlmodel import select
|
||||
|
||||
from treg.infra.db import session_maker
|
||||
from treg.models import IdempotentCall
|
||||
|
||||
async with session_maker() as db:
|
||||
return (await db.execute(select(IdempotentCall).where(
|
||||
IdempotentCall.key == key))).scalar_one_or_none()
|
||||
|
||||
|
||||
async def _set_claim(key: str, **values) -> None:
|
||||
from sqlalchemy import update
|
||||
|
||||
from treg.infra.db import session_maker
|
||||
from treg.models import IdempotentCall
|
||||
|
||||
async with session_maker() as db:
|
||||
await db.execute(update(IdempotentCall).where(IdempotentCall.key == key).values(**values))
|
||||
await db.commit()
|
||||
|
||||
|
||||
def _ago(**kw):
|
||||
from datetime import timedelta
|
||||
return datetime.now(timezone.utc).replace(tzinfo=None) - timedelta(**kw)
|
||||
|
||||
|
||||
async def _ledger(org_id: int, call_id: str, kind: str, reason: str = "", amount: int = 0) -> None:
|
||||
import uuid
|
||||
|
||||
from treg.infra.db import session_maker
|
||||
from treg.models import LedgerEntry
|
||||
|
||||
async with session_maker() as db:
|
||||
db.add(LedgerEntry(id=uuid.uuid4().hex, org_id=org_id, kind=kind, amount_micro=amount,
|
||||
call_id=call_id, endpoint_id=EP, meta={"reason": reason} if reason else {}))
|
||||
await db.commit()
|
||||
|
||||
|
||||
async def test_an_abandoned_claim_is_closed_with_a_stored_410_never_run_again(
|
||||
clients: AsyncClient, platform_on, monkeypatch):
|
||||
"""Even an owner whose visible money shows only a clean release is not proof the operation
|
||||
finished: a later child (overflow polling, an async worker) can still settle. So the key is
|
||||
never run again; a live lease answers 409, a stale one a stored 410, and no provider call or
|
||||
charge happens under it."""
|
||||
from treg.application.call import idempotency
|
||||
|
||||
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
|
||||
await _seed_answer(clients, "stuck-label", status="pending")
|
||||
await _set_claim("stuck-label", call_ref="gone-owner")
|
||||
await _ledger(org_id, "gone-owner", "reserve", amount=-100)
|
||||
await _ledger(org_id, "gone-owner", "release", reason="not_billable_502", amount=100)
|
||||
live = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "stuck-label"})
|
||||
assert live.status_code == 409, live.text
|
||||
|
||||
balance = (await clients.get(f"/orgs/{org_id}/balance")).json()["balance_micro"]
|
||||
monkeypatch.setattr(idempotency, "IDEMPOTENCY_STALE_PENDING_S", 0)
|
||||
for _ in range(2):
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "stuck-label"})
|
||||
assert r.status_code == 410, r.text
|
||||
assert r.json()["detail"]["error"] == "idempotency_outcome_unknown"
|
||||
assert "X-Treg-Original-Cost-Micro" not in r.headers, "the original cost is not known"
|
||||
assert r.headers["X-Treg-Call-Id"] == "gone-owner"
|
||||
assert (await clients.get(f"/orgs/{org_id}/balance")).json()["balance_micro"] == balance
|
||||
assert (await _claim_row("stuck-label")).call_ref == "gone-owner"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("trail", ["none", "reaped"])
|
||||
async def test_an_abandoned_claim_with_no_proof_of_finishing_fails_closed(
|
||||
clients: AsyncClient, platform_on, monkeypatch, trail):
|
||||
"""No money trail (an unmetered or pre-reserve owner), or a hold the stale-hold reaper released
|
||||
while the owner may still have been upstream: the outcome is unknown, so a stored 410 instead
|
||||
of a second upstream call."""
|
||||
from treg.application.call import idempotency
|
||||
|
||||
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
|
||||
await _seed_answer(clients, "unknown-label", status="pending")
|
||||
await _set_claim("unknown-label", call_ref="silent-owner")
|
||||
if trail == "reaped":
|
||||
await _ledger(org_id, "silent-owner", "reserve", amount=-100)
|
||||
await _ledger(org_id, "silent-owner", "release", reason="stale_hold_reaped", amount=100)
|
||||
monkeypatch.setattr(idempotency, "IDEMPOTENCY_STALE_PENDING_S", 0)
|
||||
for _ in range(2):
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "unknown-label"})
|
||||
assert r.status_code == 410, r.text
|
||||
assert r.json()["detail"]["error"] == "idempotency_outcome_unknown"
|
||||
assert "X-Treg-Original-Cost-Micro" not in r.headers, "the original cost is not known"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("task_status", ["settled", "pending"])
|
||||
async def test_an_async_child_that_settles_after_the_owner_died_is_still_a_charge(
|
||||
clients: AsyncClient, platform_on, monkeypatch, task_status):
|
||||
"""A routed owner whose early child released cleanly and whose later async child is settled
|
||||
by the worker AFTER the owner's lifetime: that late settle is a charge (410, never a takeover
|
||||
that bills the key again), and a task still pending is money in flight (409)."""
|
||||
import uuid
|
||||
|
||||
from treg.application.call import idempotency
|
||||
from treg.infra.db import session_maker
|
||||
from treg.models import AsyncTaskRecord, LedgerEntry
|
||||
|
||||
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
|
||||
await _seed_answer(clients, "async-label", status="pending")
|
||||
await _set_claim("async-label", call_ref="routed-owner", created_at=_ago(hours=2))
|
||||
stamp = _ago(hours=2)
|
||||
async with session_maker() as db:
|
||||
db.add(LedgerEntry(id=uuid.uuid4().hex, org_id=org_id, kind="release", amount_micro=100,
|
||||
call_id="routed-owner:r0", endpoint_id=EP,
|
||||
meta={"reason": "not_billable_404"}, created_at=stamp))
|
||||
db.add(AsyncTaskRecord(call_id="routed-owner:r1", org_id=org_id, provider="p", endpoint_id=EP,
|
||||
reserved_micro=500, next_check_at=stamp, status=task_status,
|
||||
created_at=stamp))
|
||||
if task_status == "settled": # the worker settles long after the owner stopped
|
||||
db.add(LedgerEntry(id=uuid.uuid4().hex, org_id=org_id, kind="settle", amount_micro=-500,
|
||||
call_id="routed-owner:r1", endpoint_id=EP, meta={}))
|
||||
await db.commit()
|
||||
monkeypatch.setattr(idempotency, "IDEMPOTENCY_STALE_PENDING_S", 60)
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "async-label"})
|
||||
if task_status == "pending":
|
||||
assert r.status_code == 409, r.text
|
||||
else:
|
||||
assert r.status_code == 410, r.text
|
||||
assert r.json()["detail"]["error"] == "idempotency_response_lost"
|
||||
assert r.json()["detail"]["charged_micro"] == 500
|
||||
|
||||
|
||||
async def test_a_renewal_between_read_and_close_keeps_the_lease(clients: AsyncClient, platform_on):
|
||||
"""The compare-and-swap includes the lease timestamp: an owner that renews after the retry
|
||||
read the stale row keeps its claim."""
|
||||
from types import SimpleNamespace
|
||||
|
||||
from treg.application.call import idempotency
|
||||
from treg.infra.db import session_maker
|
||||
|
||||
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
|
||||
await _seed_answer(clients, "race-label", status="pending")
|
||||
await _set_claim("race-label", call_ref="alive-owner", created_at=_ago(hours=1))
|
||||
await _ledger(org_id, "alive-owner", "release", reason="not_billable_502", amount=100)
|
||||
stale = await _claim_row("race-label")
|
||||
await idempotency._renew_claim_lease((stale.membership_id, stale.key, "alive-owner"))
|
||||
caller = SimpleNamespace(org_id=org_id, membership=SimpleNamespace(id=stale.membership_id))
|
||||
async with session_maker() as db:
|
||||
out = await idempotency._resolve_stale_claim(stale, caller, db)
|
||||
assert out is None
|
||||
row = await _claim_row("race-label")
|
||||
assert row.call_ref == "alive-owner" and row.status == "pending"
|
||||
|
||||
|
||||
async def test_an_abandoned_CHARGED_claim_answers_410_and_never_charges_twice(
|
||||
clients: AsyncClient, platform_on, monkeypatch):
|
||||
"""The incident: the owner was charged, then marking its claim done failed. Forgetting the claim
|
||||
would bill the same key again, so the retry gets a stored 410 naming the charge instead."""
|
||||
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
|
||||
first = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "paid-label"})
|
||||
assert first.status_code == 200, first.text
|
||||
owner = first.headers["X-Treg-Call-Id"]
|
||||
from treg.application.call import idempotency
|
||||
|
||||
await _set_claim("paid-label", status="pending", response_status=None, response_body=None,
|
||||
charged_micro=0)
|
||||
monkeypatch.setattr(idempotency, "IDEMPOTENCY_STALE_PENDING_S", 0)
|
||||
balance = (await clients.get(f"/orgs/{org_id}/balance")).json()["balance_micro"]
|
||||
|
||||
for _ in range(2):
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "paid-label"})
|
||||
assert r.status_code == 410, r.text
|
||||
detail = r.json()["detail"]
|
||||
assert detail["error"] == "idempotency_response_lost" and detail["call_id"] == owner
|
||||
assert detail["charged_micro"] > 0
|
||||
assert r.headers["X-Treg-Call-Id"] == owner
|
||||
assert r.headers["X-Treg-Original-Cost-Micro"] == str(detail["charged_micro"])
|
||||
assert (await clients.get(f"/orgs/{org_id}/balance")).json()["balance_micro"] == balance
|
||||
|
||||
|
||||
async def test_an_expired_lease_with_money_in_flight_still_answers_409(
|
||||
clients: AsyncClient, platform_on):
|
||||
"""An open hold under the owner's call id (or a child's) means the call or its async worker is
|
||||
not finished, however old the lease."""
|
||||
from treg.infra.db import session_maker
|
||||
from treg.models import Hold
|
||||
|
||||
org_id = (await clients.get("/orgs")).json()[0]["org_id"]
|
||||
await _seed_answer(clients, "busy-label", status="pending")
|
||||
await _set_claim("busy-label", call_ref="busy-owner", created_at=_ago(hours=1))
|
||||
async with session_maker() as db:
|
||||
db.add(Hold(id="busy-owner:r2", org_id=org_id, endpoint_id=EP, amount_micro=100))
|
||||
await db.commit()
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "busy-label"})
|
||||
assert r.status_code == 409, r.text
|
||||
|
||||
|
||||
async def test_a_legacy_claim_without_an_owner_is_never_taken_over(clients: AsyncClient, platform_on):
|
||||
"""Rows written before claims carried their call id have no money to check: they wait out the
|
||||
window rather than risk a second charge."""
|
||||
await _seed_answer(clients, "legacy-label", status="pending")
|
||||
await _set_claim("legacy-label", created_at=_ago(hours=1))
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "legacy-label"})
|
||||
assert r.status_code == 409, r.text
|
||||
|
||||
|
||||
async def test_a_renewal_during_closing_answers_409_not_a_phantom_410(
|
||||
clients: AsyncClient, platform_on, monkeypatch):
|
||||
"""Through the real replay path: the owner renews while the retry is deciding. The lost swap
|
||||
must not leave the retry's loaded row looking closed."""
|
||||
from treg.application.call import idempotency
|
||||
|
||||
await _seed_answer(clients, "phantom-label", status="pending")
|
||||
await _set_claim("phantom-label", call_ref="alive-owner", created_at=_ago(hours=1))
|
||||
real = idempotency._money_of
|
||||
|
||||
async def renew_meanwhile(db, org_id, call_ref, since, until):
|
||||
row = await _claim_row("phantom-label")
|
||||
await idempotency._renew_claim_lease((row.membership_id, row.key, "alive-owner"))
|
||||
return await real(db, org_id, call_ref, since, until)
|
||||
|
||||
monkeypatch.setattr(idempotency, "_money_of", renew_meanwhile)
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "phantom-label"})
|
||||
assert r.status_code == 409, r.text
|
||||
row = await _claim_row("phantom-label")
|
||||
assert row.status == "pending" and row.call_ref == "alive-owner"
|
||||
|
||||
|
||||
async def test_the_old_owner_cannot_touch_a_claim_taken_over(clients: AsyncClient, platform_on):
|
||||
"""After a takeover, the late original call's release and store are fenced out."""
|
||||
from types import SimpleNamespace
|
||||
|
||||
from treg.application.call import idempotency
|
||||
|
||||
await _seed_answer(clients, "fenced-label", status="pending")
|
||||
await _set_claim("fenced-label", call_ref="new-owner")
|
||||
row = await _claim_row("fenced-label")
|
||||
await idempotency._release_idempotent_claim((row.membership_id, row.key, "old-owner"))
|
||||
caller = SimpleNamespace(membership=SimpleNamespace(id=row.membership_id))
|
||||
await idempotency._store_idempotent(
|
||||
row.key, caller, status_code=200, body=b"{}", media_type="application/json",
|
||||
charged_micro=1, metered=True, call_ref="old-owner")
|
||||
after = await _claim_row("fenced-label")
|
||||
assert after is not None and after.status == "pending" and after.call_ref == "new-owner"
|
||||
|
||||
|
||||
async def test_marking_a_claim_done_survives_one_pool_timeout(
|
||||
clients: AsyncClient, platform_on, monkeypatch):
|
||||
"""The saturated pool that stranded claims: one timeout is retried on a fresh session."""
|
||||
from sqlalchemy.exc import TimeoutError as PoolTimeoutError
|
||||
|
||||
from treg.application.call import idempotency
|
||||
|
||||
real, failed = idempotency.session_maker, []
|
||||
|
||||
def flaky():
|
||||
if not failed:
|
||||
failed.append(1)
|
||||
raise PoolTimeoutError("QueuePool limit reached")
|
||||
return real()
|
||||
|
||||
monkeypatch.setattr(idempotency, "session_maker", flaky)
|
||||
monkeypatch.setattr(idempotency, "_POOL_RETRY_PAUSE_S", 0)
|
||||
r = await clients.get(f"/call/{EP}?aweme_id=7", headers={"Idempotency-Key": "flaky-label"})
|
||||
assert r.status_code == 200, r.text
|
||||
assert failed, "the store must have hit the simulated timeout"
|
||||
assert (await _claim_row("flaky-label")).status == "done"
|
||||
|
||||
|
||||
async def test_the_owner_renews_its_lease_and_nobody_else_does(clients: AsyncClient, platform_on):
|
||||
from treg.application.call import idempotency
|
||||
|
||||
await _seed_answer(clients, "lease-label", status="pending")
|
||||
await _set_claim("lease-label", call_ref="owner", created_at=_ago(hours=1))
|
||||
row = await _claim_row("lease-label")
|
||||
await idempotency._renew_claim_lease((row.membership_id, row.key, "stranger"))
|
||||
assert (await _claim_row("lease-label")).created_at < _ago(minutes=30)
|
||||
await idempotency._renew_claim_lease((row.membership_id, row.key, "owner"))
|
||||
assert (await _claim_row("lease-label")).created_at > _ago(minutes=1)
|
||||
|
||||
|
||||
async def test_the_sweep_leaves_OTHER_callers_rows_alone(clients: AsyncClient, platform_on):
|
||||
"""Scoped to the caller doing the work. A sweep that reached across callers would be a caller
|
||||
able to delete another's stored answers by making one call of their own."""
|
||||
|
||||
Reference in New Issue
Block a user