From a6ef60249f1ddd972759a3e10ef3341839634074 Mon Sep 17 00:00:00 2001 From: Tom Ribuot Date: Tue, 6 Oct 2026 03:28:16 +0200 Subject: [PATCH] feat: usage report transaction IDs and the usage history Kaiten now makes usage reports idempotent with a client key and keeps their history. - openapi/ is synced from kaitencloud/kaiten@ef755b4 (feat/billing) and the models regenerated; UsageReport.behavior is named UsageBehavior. - instances.report_usage(transaction_id=...): a keyed report is sent as idempotent, so it is retried like a read under the same key; a report without a key is still sent once. A malformed key raises ValueError before any request. - The usage endpoint's 409 is named by its code: TransactionIdReusedError (a ConflictError, never retried) or ThresholdExceededError. - instances.report_usage_detailed() returns UsageReportResult: usage, replayed (Idempotent-Replayed) and metadata_dropped (Kaiten-Metadata-Dropped). - An API without transaction_id refuses the field: the report is resent without it, and the client stops sending keys for its lifetime, logging once. - instances.list_usage_reports() walks the history by afterSeq, stops on a short page or at its limit, and raises PaginationError on a cursor that does not move. - instances.export_usage_reports() and export_organization_usage_reports() stream CSV or NDJSON into a binary destination through the new APIClient.stream_to, which retries only before the first byte. - unasync rewrites aread and aiter_* for the sync client. Signed-off-by: Tom Ribuot --- CHANGELOG.md | 17 + README.md | 48 ++ openapi/coverage.yaml | 5 +- openapi/openapi.yaml | 490 +++++++++++++++++- openapi/source.yaml | 6 +- scripts/generate_models.py | 1 + scripts/unasync.py | 2 + src/kaitencloud/__init__.py | 4 + src/kaitencloud/_async/_http.py | 54 +- src/kaitencloud/_async/resources/instances.py | 285 +++++++++- src/kaitencloud/_exceptions.py | 12 + src/kaitencloud/_results.py | 31 ++ src/kaitencloud/_sync/_http.py | 54 +- src/kaitencloud/_sync/resources/instances.py | 285 +++++++++- src/kaitencloud/types/__init__.py | 4 + src/kaitencloud/types/_core.py | 79 +++ tests/test_contract.py | 37 ++ tests/test_usage_reports.py | 277 ++++++++++ 18 files changed, 1651 insertions(+), 40 deletions(-) create mode 100644 src/kaitencloud/_results.py create mode 100644 tests/test_usage_reports.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 150cea4..4e13cb4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,23 @@ All notable changes to this project are documented here. The format follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html): only a major version changes the public API, a minor version adds to it, and a patch version fixes it. +## [Unreleased] + +### Added + +- `transaction_id` on `instances.report_usage()`: the API applies a keyed report at most once, + so a keyed report is retried like a read, every attempt under the same key; a report without a + key is still sent once. A malformed key raises `ValueError` before any request. +- `TransactionIdReusedError`, a `ConflictError` raised when a key was already used for a + different report; never retried. A conflict over the threshold stays `ThresholdExceededError`. +- `instances.report_usage_detailed()` returns a `UsageReportResult` with the usage, `replayed` + and `metadata_dropped`. +- Against an API without `transaction_id`, the report is resent without it and the client stops + sending keys for its lifetime, logging a warning once. +- `instances.list_usage_reports()`, `instances.export_usage_reports()` and + `instances.export_organization_usage_reports()`: the usage history, and CSV/NDJSON exports + streamed into a binary destination. + ## [1.0.0] - 2026-10-01 The first release of the Kaiten Python SDK, generated and tested against the contract of diff --git a/README.md b/README.md index 6cae9d8..9e56476 100644 --- a/README.md +++ b/README.md @@ -264,6 +264,50 @@ except ThresholdExceededError as error: `behavior="append"` (the default) folds the value into the total through the entitlement's aggregation method; `behavior="set"` replaces the total, which is also how to correct it. +A report without a key is sent once and never retried: if its response is lost, it may or may +not have been counted. Give it a `transaction_id` and the API applies it at most once (per +instance, entitlement and key, within 35 days by default), so the SDK retries it like a read, +every attempt under the same key: + +```python +from kaitencloud import TransactionIdReusedError + +try: + result = client.instances.report_usage_detailed( + "acme-production", + "tokens", + 1200, + transaction_id="llm-call:9f2c:tokens", # a UUID, or the business event id plus the meter + ) +except TransactionIdReusedError as error: + # The key was already used for a different report: a bug in how keys are made, never + # retried. A correction is a new report under a new key. + print("Original report:", error.errors[0].value) +else: + print(result.usage.value, result.replayed, result.metadata_dropped) +``` + +`result.replayed` says the API had already counted this report and `result.usage` is its +original answer; `result.metadata_dropped` that the metadata was above 4 KiB and not stored. An +API older than `transaction_id` refuses the field: the report is then resent without it, and the +client stops sending keys for its lifetime, logging a warning once. + +The usage history lists every accepted report, with the counter before and after it and the +limit it was gated on. Exports stream into a binary file as they arrive: + +```python +reports = client.instances.list_usage_reports("acme-production", "tokens", from_="2026-10-01T00:00:00Z") + +with open("tokens.csv", "wb") as destination: + client.instances.export_usage_reports("acme-production", "tokens", destination) + +# The whole organization, 31 days at a time; instance_id reaches a deleted instance. +with open("september.ndjson", "wb") as destination: + client.instances.export_organization_usage_reports( + destination, from_="2026-09-01T00:00:00Z", to="2026-10-01T00:00:00Z", format="json" + ) +``` + `kaitencloud.usage` computes the same boundaries the server enforces, to render a meter or check a report before sending it: @@ -665,6 +709,10 @@ all 10 of its operations. The async clients have the same methods. | `client.instances.list_usage()` | `GET /instances/{instanceSlug}/entitlements/usage` | | `client.instances.get_usage()` | `GET /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage` | | `client.instances.report_usage()` | `POST /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage` | +| `client.instances.report_usage_detailed()` | `POST /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage` | +| `client.instances.list_usage_reports()` | `GET /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage/reports` | +| `client.instances.export_usage_reports()` | `GET /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage/reports/export` | +| `client.instances.export_organization_usage_reports()` | `GET /usage/reports/export` | | `client.instances.get_integration()` | `GET /instances/{instanceSlug}/integrations/{integrationName}` | | `client.instances.create_integration()` | `POST /instances/{instanceSlug}/integrations/{integrationName}` | | `client.instances.update_integration()` | `PUT /instances/{instanceSlug}/integrations/{integrationName}` | diff --git a/openapi/coverage.yaml b/openapi/coverage.yaml index 7a099de..109bdff 100644 --- a/openapi/coverage.yaml +++ b/openapi/coverage.yaml @@ -82,7 +82,10 @@ core: getAuditTrails: instances.list_audit_trails getEntitlementsUsageMetrics: instances.list_usage getEntitlementUsageMetrics: instances.get_usage - reportEntitlementUsageMetric: instances.report_usage + reportEntitlementUsageMetric: [instances.report_usage, instances.report_usage_detailed] + listUsageReports: instances.list_usage_reports + exportUsageReports: instances.export_usage_reports + exportOrganizationUsageReports: instances.export_organization_usage_reports get-instance-integration: instances.get_integration create-instance-integration: instances.create_integration update-instance-integration: instances.update_integration diff --git a/openapi/openapi.yaml b/openapi/openapi.yaml index e09e19c..0d21213 100644 --- a/openapi/openapi.yaml +++ b/openapi/openapi.yaml @@ -2720,8 +2720,13 @@ components: type: string metadata: additionalProperties: {} - description: Optional metadata for the usage report + description: "Optional metadata for the usage report, a JSON object stored with it in the usage history when its compact encoding is at most 4 KiB. Above that it is not stored, the report is still counted, and the response carries Kaiten-Metadata-Dropped: too_large. It must contain no personal data: anyone who can read the organization's instances can read it, for as long as the usage history is kept. Numbers are read as 64-bit floats, so send large identifiers as strings." type: object + transactionId: + description: "Optional idempotency key, 1 to 128 characters of [A-Za-z0-9._:-], matched exactly and case-sensitively. A report sent again with the same key and the same behavior and value within KAITEN_USAGE_IDEMPOTENCY_WINDOW (35 days by default) is applied once: the retry answers 200 with the original response and the Idempotent-Replayed header, and changes nothing. The same key with another behavior or value answers 409 ReportEntitlementUsageMetric.TransactionIdReused. A rejected report does not consume its key. Scoped to the instance and entitlement: one business event may feed two meters under one key." + examples: + - llm-call-9f2c:tokens + type: string value: description: Reported entitlement value, discriminated by the 'type' field. Usage reporting accepts the number variant only. discriminator: @@ -3150,6 +3155,135 @@ components: - createdAt - createdBy type: object + UsageReport: + additionalProperties: false + properties: + aggregationMethod: + description: The entitlement's aggregation method when the report was accepted + examples: + - SUM + type: string + behavior: + description: append adds the value through the aggregation method; set overwrites the counter + enum: + - append + - set + type: string + delta: + description: valueAfter minus valueBefore; negative for a set that lowered the counter + examples: + - "600" + type: string + entitlementId: + description: The entitlement the report was made for. It may since have been deleted. + format: uuid + type: string + eventCountAfter: + description: Reports counted in the window after this one + examples: + - 1 + format: int32 + type: integer + instanceId: + description: The instance the report was made for. It may since have been deleted. + format: uuid + type: string + licenseId: + description: The instance's licence when the report was accepted + format: uuid + type: string + limitValue: + description: The limit in force when the report was accepted. Null when unlimited. + examples: + - "1000" + type: string + overageDelta: + description: "How much the usage above limitValue moved: max(0, valueAfter - limitValue) - max(0, valueBefore - limitValue). 0 when unlimited." + examples: + - "0" + type: string + overagePercent: + description: The overage allowed above limitValue, in percent, when the report was accepted. Null when unlimited. + examples: + - 50 + format: int32 + type: integer + properties: + additionalProperties: {} + description: The report's metadata, when it was stored + type: object + reportSeq: + description: Position of the report among the pair's accepted reports, from 1, without gaps + examples: + - 42 + format: int64 + type: integer + reportedAt: + description: When the server accepted the report (UTC, millisecond precision) + examples: + - "2026-10-05T08:00:00.000Z" + format: date-time + type: string + reportedValue: + description: "The value as sent: a delta for append, an absolute value for set" + examples: + - "600" + type: string + transactionId: + description: The report's idempotency key, when it was sent with one + examples: + - llm-call-9f2c:tokens + type: string + valueAfter: + description: The counter after the report + examples: + - "600" + type: string + valueBefore: + description: The counter before the report, after any window reset + examples: + - "0" + type: string + windowEnd: + description: End of the usage window the report counted in (exclusive). Null for a lifetime entitlement. + format: date-time + type: string + windowStart: + description: Start of the usage window the report counted in (inclusive). Null for a lifetime entitlement. + format: date-time + type: string + required: + - instanceId + - entitlementId + - reportSeq + - reportedAt + - behavior + - aggregationMethod + - reportedValue + - valueBefore + - valueAfter + - delta + - overageDelta + - eventCountAfter + - licenseId + type: object + UsageReportPage: + additionalProperties: false + properties: + items: + description: The reports, in reportSeq order + items: + $ref: "#/components/schemas/UsageReport" + type: array + nextAfterSeq: + description: Pass as afterSeq to read the next page. Absent on the last page. + examples: + - 100 + format: int64 + type: integer + required: + - items + type: object User: additionalProperties: false properties: @@ -6850,7 +6984,7 @@ paths: tags: - instances post: - description: Report a usage metric for a specific entitlement in a given instance. This endpoint allows you to report the usage of an entitlement, including optional metadata and a timestamp. + description: "Report a usage metric for a specific entitlement in a given instance, with optional metadata. The server dates every report on receipt; the request carries no timestamp. Send a transactionId to make retries safe: without one, a report sent twice counts twice." operationId: reportEntitlementUsageMetric parameters: - description: Instance slug @@ -6884,6 +7018,15 @@ paths: schema: $ref: "#/components/schemas/EntitlementUsage" description: OK + headers: + Idempotent-Replayed: + schema: + description: "true when the report replays an earlier one sent with the same transactionId: the body is that report's original response and nothing was counted again" + type: string + Kaiten-Metadata-Dropped: + schema: + description: too_large when the report's metadata was above 4 KiB and was not stored; the report itself was counted + type: string "400": content: application/problem+json: @@ -6926,12 +7069,238 @@ paths: schema: $ref: "#/components/schemas/Problem" description: Internal Server Error + "503": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Service Unavailable security: - bearerAuth: - write:instances summary: Report entitlement usage metric for an instance tags: - instances + /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage/reports: + get: + description: "The usage history of one instance and entitlement: every accepted report in the range, with the counter before and after it and the limit in force, in reportSeq order, paged with afterSeq. Decimals are strings. Reports are kept for the organization's usage history retention." + operationId: listUsageReports + parameters: + - description: Instance slug + in: path + name: instanceSlug + required: true + schema: + description: Instance slug + examples: + - instance-slug + type: string + - description: Entitlement slug + in: path + name: entitlementSlug + required: true + schema: + description: Entitlement slug + examples: + - entitlement-slug + type: string + - description: Start of the range, inclusive (RFC 3339). Defaults to 30 days before to, moved up to the start of the organization's usage history when that is later. An explicit from before it answers 422 ListUsageReports.OutsideRetention. + explode: false + in: query + name: from + schema: + description: Start of the range, inclusive (RFC 3339). Defaults to 30 days before to, moved up to the start of the organization's usage history when that is later. An explicit from before it answers 422 ListUsageReports.OutsideRetention. + format: date-time + type: string + - description: End of the range, exclusive (RFC 3339). Defaults to now. + explode: false + in: query + name: to + schema: + description: End of the range, exclusive (RFC 3339). Defaults to now. + format: date-time + type: string + - description: "Return the reports after this reportSeq: the nextAfterSeq of the previous page" + explode: false + in: query + name: afterSeq + schema: + description: "Return the reports after this reportSeq: the nextAfterSeq of the previous page" + format: int64 + minimum: 0 + type: integer + - description: Maximum number of reports to return (default 100, max 500) + explode: false + in: query + name: limit + schema: + description: Maximum number of reports to return (default 100, max 500) + format: int32 + maximum: 500 + minimum: 1 + type: integer + - description: Only the report sent with this idempotency key + explode: false + in: query + name: transactionId + schema: + description: Only the report sent with this idempotency key + maxLength: 128 + type: string + responses: + "200": + content: + application/json: + schema: + $ref: "#/components/schemas/UsageReportPage" + description: OK + "400": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Bad Request + "401": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Unauthorized + "403": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Forbidden + "404": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Not Found + "422": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Unprocessable Entity + "500": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Internal Server Error + security: + - bearerAuth: + - read:instances + summary: List an entitlement's usage reports for an instance + tags: + - instances + /instances/{instanceSlug}/entitlements/{entitlementSlug}/usage/reports/export: + get: + description: Streams every report of one instance and entitlement in the range as CSV or NDJSON, in reportSeq order. Decimals are written exactly as the journal stores them. + operationId: exportUsageReports + parameters: + - description: Instance slug + in: path + name: instanceSlug + required: true + schema: + description: Instance slug + examples: + - instance-slug + type: string + - description: Entitlement slug + in: path + name: entitlementSlug + required: true + schema: + description: Entitlement slug + examples: + - entitlement-slug + type: string + - description: Start of the range, inclusive (RFC 3339). Defaults to 30 days before to, moved up to the start of the organization's usage history when that is later. An explicit from before it answers 422 ExportUsageReports.OutsideRetention. + explode: false + in: query + name: from + schema: + description: Start of the range, inclusive (RFC 3339). Defaults to 30 days before to, moved up to the start of the organization's usage history when that is later. An explicit from before it answers 422 ExportUsageReports.OutsideRetention. + format: date-time + type: string + - description: End of the range, exclusive (RFC 3339). Defaults to now. At most 366 days after from. + explode: false + in: query + name: to + schema: + description: End of the range, exclusive (RFC 3339). Defaults to now. At most 366 days after from. + format: date-time + type: string + - description: "csv (the default): RFC 4180 with a header row. json: NDJSON, one report per line, in the shape listUsageReports returns. Anything else answers 422 ExportUsageReports.InvalidFormat." + explode: false + in: query + name: format + schema: + description: "csv (the default): RFC 4180 with a header row. json: NDJSON, one report per line, in the shape listUsageReports returns. Anything else answers 422 ExportUsageReports.InvalidFormat." + examples: + - csv + type: string + responses: + "200": + content: + application/x-ndjson: + schema: + type: string + text/csv: + schema: + type: string + description: "The reports, as an attachment: CSV with a header row for format=csv, one JSON object per line for format=json" + headers: + Content-Disposition: + description: attachment, with a file name naming the range + schema: + type: string + "400": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Bad Request + "401": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Unauthorized + "403": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Forbidden + "404": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Not Found + "422": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Unprocessable Entity + "500": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Internal Server Error + security: + - bearerAuth: + - read:instances + summary: Export an entitlement's usage reports for an instance + tags: + - instances /instances/{instanceSlug}/integrations/{integrationName}: delete: description: Delete a single integration for an instance @@ -9939,6 +10308,123 @@ paths: summary: Revoke a token for a service account tags: - service-accounts + /usage/reports/export: + get: + description: Streams every usage report of the organization in the range, across instances and entitlements, as CSV or NDJSON, ordered by reportedAt. Filters narrow it to one instance or one entitlement, including deleted ones by ID. Use it to keep the usage history before deleting the organization. + operationId: exportOrganizationUsageReports + parameters: + - description: Start of the range, inclusive (RFC 3339). Defaults to 30 days before to, moved up to the start of the organization's usage history when that is later. An explicit from before it answers 422 ExportUsageReports.OutsideRetention. + explode: false + in: query + name: from + schema: + description: Start of the range, inclusive (RFC 3339). Defaults to 30 days before to, moved up to the start of the organization's usage history when that is later. An explicit from before it answers 422 ExportUsageReports.OutsideRetention. + format: date-time + type: string + - description: End of the range, exclusive (RFC 3339). Defaults to now. At most 31 days after from. + explode: false + in: query + name: to + schema: + description: End of the range, exclusive (RFC 3339). Defaults to now. At most 31 days after from. + format: date-time + type: string + - description: "csv (the default): RFC 4180 with a header row. json: NDJSON, one report per line. Anything else answers 422 ExportUsageReports.InvalidFormat." + explode: false + in: query + name: format + schema: + description: "csv (the default): RFC 4180 with a header row. json: NDJSON, one report per line. Anything else answers 422 ExportUsageReports.InvalidFormat." + examples: + - csv + type: string + - description: Only this instance's reports + explode: false + in: query + name: instanceSlug + schema: + description: Only this instance's reports + type: string + - description: Only this instance's reports. Reaches a deleted instance, whose reports are kept. + explode: false + in: query + name: instanceId + schema: + description: Only this instance's reports. Reaches a deleted instance, whose reports are kept. + format: uuid + type: string + - description: Only this entitlement's reports + explode: false + in: query + name: entitlementSlug + schema: + description: Only this entitlement's reports + type: string + - description: Only this entitlement's reports. Reaches a deleted entitlement, whose reports are kept. + explode: false + in: query + name: entitlementId + schema: + description: Only this entitlement's reports. Reaches a deleted entitlement, whose reports are kept. + format: uuid + type: string + responses: + "200": + content: + application/x-ndjson: + schema: + type: string + text/csv: + schema: + type: string + description: "The reports, as an attachment: CSV with a header row for format=csv, one JSON object per line for format=json" + headers: + Content-Disposition: + description: attachment, with a file name naming the range + schema: + type: string + "400": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Bad Request + "401": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Unauthorized + "403": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Forbidden + "404": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Not Found + "422": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Unprocessable Entity + "500": + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" + description: Internal Server Error + security: + - bearerAuth: + - read:instances + summary: Export the organization's usage reports + tags: + - instances /v1/notification-preferences: get: description: "Every notifiable event, with this user's effective value per channel: their own choice where they made one, the catalogue default otherwise. Events absent from the catalogue of a deployment never appear, which is what keeps high-volume system events out of the feed." diff --git a/openapi/source.yaml b/openapi/source.yaml index 3f7c144..11f7acb 100644 --- a/openapi/source.yaml +++ b/openapi/source.yaml @@ -1,9 +1,9 @@ # Where the two snapshots in this directory were copied from. Rewritten by # `task sync:openapi`; never edit the snapshots themselves by hand. repository: kaitencloud/kaiten -ref: origin/main -commit: fa7f554fb6464f52c2a3f1ed487f4c7383c8125d -synced_at: "2026-09-30" +ref: feat/billing +commit: ef755b4213bb12e67aad7c50714600b8593aa44c +synced_at: "2026-10-06" files: openapi.yaml: app/openapi.yaml platform-openapi.yaml: app/platform-openapi.yaml diff --git a/scripts/generate_models.py b/scripts/generate_models.py index 1ea111c..110093b 100644 --- a/scripts/generate_models.py +++ b/scripts/generate_models.py @@ -96,6 +96,7 @@ ("License", "type"): "LicenseType", ("LicenseEntitlement", "entitlementType"): "LicenseEntitlementType", ("MetadataField", "resourceType"): "MetadataResourceType", + ("UsageReport", "behavior"): "UsageBehavior", } # Attribute names that would shadow a pydantic.BaseModel member, keyed by (schema, property). diff --git a/scripts/unasync.py b/scripts/unasync.py index cb0c32b..61e0c9e 100644 --- a/scripts/unasync.py +++ b/scripts/unasync.py @@ -41,6 +41,8 @@ (re.compile(r"\b__aenter__\b"), "__enter__"), (re.compile(r"\b__aexit__\b"), "__exit__"), (re.compile(r"\baclose\b"), "close"), + (re.compile(r"\baread\b"), "read"), + (re.compile(r"\baiter_(bytes|text|lines|raw)\b"), r"iter_\1"), (re.compile(r"\basync_(\w+)"), r"\1"), (re.compile(r"\b_Async(?=[A-Z])"), "_"), (re.compile(r"\bAsync(?=[A-Z])"), ""), diff --git a/src/kaitencloud/__init__.py b/src/kaitencloud/__init__.py index c2db8fe..9ccdd8e 100644 --- a/src/kaitencloud/__init__.py +++ b/src/kaitencloud/__init__.py @@ -37,11 +37,13 @@ RateLimitError, ServiceUnavailableError, ThresholdExceededError, + TransactionIdReusedError, UnprocessableEntityError, WebhookPayloadError, WebhookVerificationError, ) from ._models import KaitenModel +from ._results import UsageReportResult from ._sync.client import KaitenClient, KaitenPlatformClient from ._types import NOT_GIVEN, NotGiven from ._version import __version__ @@ -77,7 +79,9 @@ "RateLimitError", "ServiceUnavailableError", "ThresholdExceededError", + "TransactionIdReusedError", "UnprocessableEntityError", + "UsageReportResult", "WebhookPayloadError", "WebhookVerificationError", "__version__", diff --git a/src/kaitencloud/_async/_http.py b/src/kaitencloud/_async/_http.py index b96196c..a43e850 100644 --- a/src/kaitencloud/_async/_http.py +++ b/src/kaitencloud/_async/_http.py @@ -7,7 +7,7 @@ import logging from collections.abc import AsyncIterator, Mapping -from typing import Any, TypeVar +from typing import Any, BinaryIO, TypeVar import httpx from pydantic import BaseModel, ValidationError @@ -64,6 +64,9 @@ def __init__( self.http_client = http_client or httpx.AsyncClient( timeout=DEFAULT_TIMEOUT if isinstance(timeout, NotGiven) else timeout ) + # Set once the API refuses usage report transaction IDs: from then on, reports go out + # without them (and so without retries). + self.transaction_ids_unsupported = False def copy( self, @@ -138,6 +141,55 @@ async def request( continue raise make_status_error(response, overrides=errors) + async def stream_to( + self, + path: str, + destination: BinaryIO, + *, + params: Mapping[str, Any] | None = None, + ) -> int: + """GET ``path`` and write its body to ``destination`` as it arrives; return its size. + + The body is never held whole in memory. A failure before the body starts is retried + like any idempotent read; once bytes have been written, an interruption is raised, not + retried, since ``destination`` already holds part of the body. + """ + url = self.base_url + path.lstrip("/") + attempt = 0 + while True: + request = self._build( + "GET", url, params=_query(params), body=None, headers=await self._headers() + ) + try: + response = await self.http_client.send(request, stream=True) + except httpx.TransportError as error: + if attempt < self.max_retries: + await self._back_off(request, attempt, repr(error)) + attempt += 1 + continue + if isinstance(error, httpx.TimeoutException): + raise APITimeoutError(request=request) from error + raise APIConnectionError(request=request) from error + + try: + if not response.is_success: + await response.aread() + status = response.status_code + if attempt < self.max_retries and (status in (408, 429) or status >= 500): + await self._back_off( + request, attempt, f"HTTP {status}", response.headers.get("retry-after") + ) + attempt += 1 + continue + raise make_status_error(response) + written = 0 + async for chunk in response.aiter_bytes(): + destination.write(chunk) + written += len(chunk) + return written + finally: + await response.aclose() + async def get( self, path: str, model: type[ModelT], *, params: Mapping[str, Any] | None = None ) -> ModelT: diff --git a/src/kaitencloud/_async/resources/instances.py b/src/kaitencloud/_async/resources/instances.py index d7b6906..0712eeb 100644 --- a/src/kaitencloud/_async/resources/instances.py +++ b/src/kaitencloud/_async/resources/instances.py @@ -4,11 +4,19 @@ from __future__ import annotations import builtins +import logging +import re from collections.abc import Mapping from datetime import datetime -from typing import Any +from typing import Any, BinaryIO, Literal -from ..._exceptions import ThresholdExceededError +from ..._exceptions import ( + APIStatusError, + PaginationError, + ThresholdExceededError, + TransactionIdReusedError, +) +from ..._results import UsageReportResult from ..._utils import ( api_path, compact, @@ -25,11 +33,24 @@ InstanceStatus, IntegrationParam, UsageBehavior, + UsageReport, + UsageReportPage, ) from ._base import AsyncResource __all__ = ["AsyncInstances"] +log = logging.getLogger("kaitencloud") + +# The format the API accepts for a usage report key. +_TRANSACTION_ID = re.compile(r"[A-Za-z0-9._:-]{1,128}") +_TRANSACTION_ID_REUSED = "ReportEntitlementUsageMetric.TransactionIdReused" +_INVALID_TRANSACTION_ID = "ReportEntitlementUsageMetric.InvalidTransactionId" +# The largest page the usage history endpoint serves. +_USAGE_HISTORY_PAGE = 500 + +UsageExportFormat = Literal["csv", "json"] + class AsyncInstances(AsyncResource): """Instances: deployments of your product for a customer, under a license. @@ -221,6 +242,7 @@ async def report_usage( *, behavior: UsageBehavior = "append", metadata: Mapping[str, Any] | None = None, + transaction_id: str | None = None, ) -> EntitlementUsage: """Report usage of a NUMBER entitlement, and return the usage it results in. @@ -231,30 +253,250 @@ async def report_usage( through the entitlement's aggregation method (a sum, by default). With ``"set"`` it replaces the total -- which is also how to correct it downwards. behavior: ``"append"`` (the default) or ``"set"``. - metadata: Free-form metadata recorded with the report. + metadata: Free-form metadata stored with the report in the usage history, when its + compact JSON encoding is at most 4 KiB. It must contain no personal data. + transaction_id: An idempotency key. The API applies a report at most once per key, + instance and entitlement within its idempotency window (35 days by default), + and answers a repeat with the original result. Use a UUID, or a business event + id plus the meter (``"llm-call:9f2c:tokens"``): 1 to 128 characters of letters, + digits, ``.``, ``_``, ``:`` and ``-``. Raises: ThresholdExceededError: The report would take usage past the maximum the license allows -- the granted value plus its overage allowance. Nothing was recorded. + TransactionIdReusedError: The key was already used for a different report. + ValueError: ``transaction_id`` is malformed; nothing was sent. - Never retried automatically: without an idempotency key on the wire, a report retried + Retried, like a read, only with a ``transaction_id``: without one, a report retried after a lost response would be counted twice. """ + result = await self.report_usage_detailed( + instance_slug, + entitlement_slug, + value, + behavior=behavior, + metadata=metadata, + transaction_id=transaction_id, + ) + return result.usage + + async def report_usage_detailed( + self, + instance_slug: str, + entitlement_slug: str, + value: float, + *, + behavior: UsageBehavior = "append", + metadata: Mapping[str, Any] | None = None, + transaction_id: str | None = None, + ) -> UsageReportResult: + """:meth:`report_usage`, also returning what the API said about the report. + + :attr:`UsageReportResult.replayed` says the report repeated one the API had already + counted, and :attr:`UsageReportResult.metadata_dropped` that its metadata was too + large to store. + + An API older than ``transaction_id`` refuses the field. The report is then sent again + without it, and this client stops sending keys (and so retrying reports) for its + lifetime, logging that once. + """ if isinstance(value, bool) or not isinstance(value, (int, float)): raise TypeError(f"usage must be a number, got {type(value).__name__}") - body = compact( - { - "value": {"type": "number", "value": value}, - "behavior": behavior, - "metadata": dict(metadata) if metadata is not None else None, - } + if transaction_id is not None and not _TRANSACTION_ID.fullmatch(transaction_id): + raise ValueError( + f"transaction_id {transaction_id!r} must be 1 to 128 characters of letters, " + "digits, '.', '_', ':' and '-'" + ) + key = None if self._api.transaction_ids_unsupported else transaction_id + path = api_path("instances", instance_slug, "entitlements", entitlement_slug, "usage") + + def body(with_key: str | None) -> dict[str, Any]: + return compact( + { + "value": {"type": "number", "value": value}, + "behavior": behavior, + "metadata": dict(metadata) if metadata is not None else None, + "transactionId": with_key, + } + ) + + try: + response = await self._send_report(path, body(key), keyed=key is not None) + except APIStatusError as error: + if key is None or not _refuses_transaction_id(error): + raise + self._api.transaction_ids_unsupported = True + log.warning( + "The Kaiten API does not support usage report transaction IDs: reports are sent " + "without them and are not retried for the rest of this client's life. Upgrade " + "Kaiten to make retries safe." + ) + response = await self._send_report(path, body(None), keyed=False) + + return UsageReportResult( + usage=self._api.validate(EntitlementUsage, self._api.decode(response), response), + replayed=response.headers.get("idempotent-replayed") == "true", + metadata_dropped=bool(response.headers.get("kaiten-metadata-dropped")), ) - return await self._api.send( - "POST", - api_path("instances", instance_slug, "entitlements", entitlement_slug, "usage"), - EntitlementUsage, - body=body, - errors={409: ThresholdExceededError}, + + async def _send_report(self, path: str, body: dict[str, Any], *, keyed: bool) -> Any: + try: + return await self._api.request( + "POST", path, body=body, idempotent=keyed, errors={409: ThresholdExceededError} + ) + except ThresholdExceededError as error: + # The usage endpoint answers 409 for two reasons; only the code tells them apart. + if error.code == _TRANSACTION_ID_REUSED: + raise TransactionIdReusedError( + str(error), response=error.response, problem=error.problem + ) from None + raise + + async def list_usage_reports( + self, + instance_slug: str, + entitlement_slug: str, + *, + from_: datetime | str | None = None, + to: datetime | str | None = None, + transaction_id: str | None = None, + limit: int | None = None, + ) -> builtins.list[UsageReport]: + """Return the usage history of one entitlement for one instance, oldest first. + + Every report the API accepted in the range, with the counter before and after it and + the limit it was gated on. Decimals are strings, exact to the digit. Pages are walked + until the range is exhausted or ``limit`` reports are held. + + Args: + instance_slug: The instance. + entitlement_slug: The entitlement. + from_: Start of the range, inclusive. Defaults to 30 days before ``to``; a value + before the start of the organization's usage history is refused with the code + ``ListUsageReports.OutsideRetention``. + to: End of the range, exclusive. Defaults to now. + transaction_id: Only the report sent with this key. + limit: The most reports to return. + """ + if limit is not None and limit < 1: + raise ValueError(f"limit must be at least 1, got {limit}") + path = api_path( + "instances", instance_slug, "entitlements", entitlement_slug, "usage", "reports" + ) + reports: builtins.list[UsageReport] = [] + after_seq: int | None = None + while True: + wanted = _USAGE_HISTORY_PAGE if limit is None else limit - len(reports) + size = min(_USAGE_HISTORY_PAGE, wanted) + page = await self._api.get( + path, + UsageReportPage, + params={ + "from": optional_datetime(from_), + "to": optional_datetime(to), + "transactionId": transaction_id, + "afterSeq": after_seq, + "limit": size, + }, + ) + items = page.items or [] + reports.extend(items) + # The API names a next page only after a full one; a short page is the last. + if ( + page.next_after_seq is None + or len(items) < size + or (limit is not None and len(reports) >= limit) + ): + return reports + if after_seq is not None and page.next_after_seq <= after_seq: + raise PaginationError( + f"GET {path}: nextAfterSeq {page.next_after_seq} does not move past " + f"{after_seq}; refusing to return a partial history as a complete one" + ) + after_seq = page.next_after_seq + + async def export_usage_reports( + self, + instance_slug: str, + entitlement_slug: str, + destination: BinaryIO, + *, + from_: datetime | str | None = None, + to: datetime | str | None = None, + format: UsageExportFormat = "csv", + ) -> int: + """Stream the usage history of one entitlement for one instance into ``destination``. + + CSV (with a header row) or NDJSON (one report per line), written as it arrives, so the + export is never held whole in memory. At most 366 days per call. Returns the number of + bytes written. + + Args: + instance_slug: The instance. + entitlement_slug: The entitlement. + destination: A binary file-like object, such as ``open(path, "wb")``. + from_: Start of the range, inclusive (see :meth:`list_usage_reports`). + to: End of the range, exclusive. Defaults to now. + format: ``"csv"`` (the default) or ``"json"`` for NDJSON. + """ + return await self._api.stream_to( + api_path( + "instances", + instance_slug, + "entitlements", + entitlement_slug, + "usage", + "reports", + "export", + ), + destination, + params={ + "from": optional_datetime(from_), + "to": optional_datetime(to), + "format": format, + }, + ) + + async def export_organization_usage_reports( + self, + destination: BinaryIO, + *, + from_: datetime | str | None = None, + to: datetime | str | None = None, + format: UsageExportFormat = "csv", + instance_slug: str | None = None, + instance_id: str | None = None, + entitlement_slug: str | None = None, + entitlement_id: str | None = None, + ) -> int: + """Stream the whole organization's usage history into ``destination``. + + Every instance and entitlement, at most 31 days per call, optionally for one instance + or one entitlement. ``instance_id`` and ``entitlement_id`` reach deleted ones, whose + reports are kept. Returns the number of bytes written. + + Args: + destination: A binary file-like object, such as ``open(path, "wb")``. + from_: Start of the range, inclusive (see :meth:`list_usage_reports`). + to: End of the range, exclusive. Defaults to now. + format: ``"csv"`` (the default) or ``"json"`` for NDJSON. + instance_slug: Only this instance. + instance_id: Only this instance, by id. + entitlement_slug: Only this entitlement. + entitlement_id: Only this entitlement, by id. + """ + return await self._api.stream_to( + api_path("usage", "reports", "export"), + destination, + params={ + "from": optional_datetime(from_), + "to": optional_datetime(to), + "format": format, + "instanceSlug": instance_slug, + "instanceId": instance_id, + "entitlementSlug": entitlement_slug, + "entitlementId": entitlement_id, + }, ) async def get_integration( @@ -347,3 +589,14 @@ def _instance( } ), } + + +def _refuses_transaction_id(error: APIStatusError) -> bool: + """Whether ``error`` is an API without transaction IDs refusing the field. + + That is a validation problem located on the field, from an API that does not know the code + newer ones answer a malformed key with. The key was checked before sending. + """ + if error.status_code not in (400, 422) or error.code == _INVALID_TRANSACTION_ID: + return False + return any(detail.location == "body.transactionId" for detail in error.errors) diff --git a/src/kaitencloud/_exceptions.py b/src/kaitencloud/_exceptions.py index 23c680e..c0d0923 100644 --- a/src/kaitencloud/_exceptions.py +++ b/src/kaitencloud/_exceptions.py @@ -30,6 +30,7 @@ "RateLimitError", "ServiceUnavailableError", "ThresholdExceededError", + "TransactionIdReusedError", "UnprocessableEntityError", "WebhookPayloadError", "WebhookVerificationError", @@ -162,6 +163,17 @@ class ThresholdExceededError(ConflictError): """ +class TransactionIdReusedError(ConflictError): + """A usage report's ``transaction_id`` was already used for a different report. + + Within the API's idempotency window, a key names one report: the same key with another + behavior or value is refused, and nothing was recorded. It is a bug in how keys are made, + not a transient failure -- the same report under the same key fails the same way, so it is + never retried. Send a correction as a new report under a new key. The original report is + in ``errors[0].value``. + """ + + class UnprocessableEntityError(APIStatusError): """The API answered 422: the request is well-formed but its content is refused.""" diff --git a/src/kaitencloud/_results.py b/src/kaitencloud/_results.py new file mode 100644 index 0000000..18492bb --- /dev/null +++ b/src/kaitencloud/_results.py @@ -0,0 +1,31 @@ +# Copyright 2026 KAITEN INC +# SPDX-License-Identifier: Apache-2.0 + +"""Results that carry more than a response body.""" + +from __future__ import annotations + +from dataclasses import dataclass + +from .types import EntitlementUsage + +__all__ = ["UsageReportResult"] + + +@dataclass(frozen=True) +class UsageReportResult: + """What an accepted usage report answers. + + Attributes: + usage: The entitlement's usage after the report. + replayed: The report carried a ``transaction_id`` the API had already accepted: + ``usage`` is that report's original answer and nothing was counted again. It may + describe a usage window that has since closed, so do not read it as the current + gauge. + metadata_dropped: The metadata was above 4 KiB and was not stored. The report itself + was counted. + """ + + usage: EntitlementUsage + replayed: bool + metadata_dropped: bool diff --git a/src/kaitencloud/_sync/_http.py b/src/kaitencloud/_sync/_http.py index 7c44ffc..f9b69b3 100644 --- a/src/kaitencloud/_sync/_http.py +++ b/src/kaitencloud/_sync/_http.py @@ -9,7 +9,7 @@ import logging from collections.abc import Iterator, Mapping -from typing import Any, TypeVar +from typing import Any, BinaryIO, TypeVar import httpx from pydantic import BaseModel, ValidationError @@ -66,6 +66,9 @@ def __init__( self.http_client = http_client or httpx.Client( timeout=DEFAULT_TIMEOUT if isinstance(timeout, NotGiven) else timeout ) + # Set once the API refuses usage report transaction IDs: from then on, reports go out + # without them (and so without retries). + self.transaction_ids_unsupported = False def copy( self, @@ -140,6 +143,55 @@ def request( continue raise make_status_error(response, overrides=errors) + def stream_to( + self, + path: str, + destination: BinaryIO, + *, + params: Mapping[str, Any] | None = None, + ) -> int: + """GET ``path`` and write its body to ``destination`` as it arrives; return its size. + + The body is never held whole in memory. A failure before the body starts is retried + like any idempotent read; once bytes have been written, an interruption is raised, not + retried, since ``destination`` already holds part of the body. + """ + url = self.base_url + path.lstrip("/") + attempt = 0 + while True: + request = self._build( + "GET", url, params=_query(params), body=None, headers=self._headers() + ) + try: + response = self.http_client.send(request, stream=True) + except httpx.TransportError as error: + if attempt < self.max_retries: + self._back_off(request, attempt, repr(error)) + attempt += 1 + continue + if isinstance(error, httpx.TimeoutException): + raise APITimeoutError(request=request) from error + raise APIConnectionError(request=request) from error + + try: + if not response.is_success: + response.read() + status = response.status_code + if attempt < self.max_retries and (status in (408, 429) or status >= 500): + self._back_off( + request, attempt, f"HTTP {status}", response.headers.get("retry-after") + ) + attempt += 1 + continue + raise make_status_error(response) + written = 0 + for chunk in response.iter_bytes(): + destination.write(chunk) + written += len(chunk) + return written + finally: + response.close() + def get( self, path: str, model: type[ModelT], *, params: Mapping[str, Any] | None = None ) -> ModelT: diff --git a/src/kaitencloud/_sync/resources/instances.py b/src/kaitencloud/_sync/resources/instances.py index ecae47d..ba67b38 100644 --- a/src/kaitencloud/_sync/resources/instances.py +++ b/src/kaitencloud/_sync/resources/instances.py @@ -6,11 +6,19 @@ from __future__ import annotations import builtins +import logging +import re from collections.abc import Mapping from datetime import datetime -from typing import Any +from typing import Any, BinaryIO, Literal -from ..._exceptions import ThresholdExceededError +from ..._exceptions import ( + APIStatusError, + PaginationError, + ThresholdExceededError, + TransactionIdReusedError, +) +from ..._results import UsageReportResult from ..._utils import ( api_path, compact, @@ -27,11 +35,24 @@ InstanceStatus, IntegrationParam, UsageBehavior, + UsageReport, + UsageReportPage, ) from ._base import Resource __all__ = ["Instances"] +log = logging.getLogger("kaitencloud") + +# The format the API accepts for a usage report key. +_TRANSACTION_ID = re.compile(r"[A-Za-z0-9._:-]{1,128}") +_TRANSACTION_ID_REUSED = "ReportEntitlementUsageMetric.TransactionIdReused" +_INVALID_TRANSACTION_ID = "ReportEntitlementUsageMetric.InvalidTransactionId" +# The largest page the usage history endpoint serves. +_USAGE_HISTORY_PAGE = 500 + +UsageExportFormat = Literal["csv", "json"] + class Instances(Resource): """Instances: deployments of your product for a customer, under a license. @@ -223,6 +244,7 @@ def report_usage( *, behavior: UsageBehavior = "append", metadata: Mapping[str, Any] | None = None, + transaction_id: str | None = None, ) -> EntitlementUsage: """Report usage of a NUMBER entitlement, and return the usage it results in. @@ -233,30 +255,250 @@ def report_usage( through the entitlement's aggregation method (a sum, by default). With ``"set"`` it replaces the total -- which is also how to correct it downwards. behavior: ``"append"`` (the default) or ``"set"``. - metadata: Free-form metadata recorded with the report. + metadata: Free-form metadata stored with the report in the usage history, when its + compact JSON encoding is at most 4 KiB. It must contain no personal data. + transaction_id: An idempotency key. The API applies a report at most once per key, + instance and entitlement within its idempotency window (35 days by default), + and answers a repeat with the original result. Use a UUID, or a business event + id plus the meter (``"llm-call:9f2c:tokens"``): 1 to 128 characters of letters, + digits, ``.``, ``_``, ``:`` and ``-``. Raises: ThresholdExceededError: The report would take usage past the maximum the license allows -- the granted value plus its overage allowance. Nothing was recorded. + TransactionIdReusedError: The key was already used for a different report. + ValueError: ``transaction_id`` is malformed; nothing was sent. - Never retried automatically: without an idempotency key on the wire, a report retried + Retried, like a read, only with a ``transaction_id``: without one, a report retried after a lost response would be counted twice. """ + result = self.report_usage_detailed( + instance_slug, + entitlement_slug, + value, + behavior=behavior, + metadata=metadata, + transaction_id=transaction_id, + ) + return result.usage + + def report_usage_detailed( + self, + instance_slug: str, + entitlement_slug: str, + value: float, + *, + behavior: UsageBehavior = "append", + metadata: Mapping[str, Any] | None = None, + transaction_id: str | None = None, + ) -> UsageReportResult: + """:meth:`report_usage`, also returning what the API said about the report. + + :attr:`UsageReportResult.replayed` says the report repeated one the API had already + counted, and :attr:`UsageReportResult.metadata_dropped` that its metadata was too + large to store. + + An API older than ``transaction_id`` refuses the field. The report is then sent again + without it, and this client stops sending keys (and so retrying reports) for its + lifetime, logging that once. + """ if isinstance(value, bool) or not isinstance(value, (int, float)): raise TypeError(f"usage must be a number, got {type(value).__name__}") - body = compact( - { - "value": {"type": "number", "value": value}, - "behavior": behavior, - "metadata": dict(metadata) if metadata is not None else None, - } + if transaction_id is not None and not _TRANSACTION_ID.fullmatch(transaction_id): + raise ValueError( + f"transaction_id {transaction_id!r} must be 1 to 128 characters of letters, " + "digits, '.', '_', ':' and '-'" + ) + key = None if self._api.transaction_ids_unsupported else transaction_id + path = api_path("instances", instance_slug, "entitlements", entitlement_slug, "usage") + + def body(with_key: str | None) -> dict[str, Any]: + return compact( + { + "value": {"type": "number", "value": value}, + "behavior": behavior, + "metadata": dict(metadata) if metadata is not None else None, + "transactionId": with_key, + } + ) + + try: + response = self._send_report(path, body(key), keyed=key is not None) + except APIStatusError as error: + if key is None or not _refuses_transaction_id(error): + raise + self._api.transaction_ids_unsupported = True + log.warning( + "The Kaiten API does not support usage report transaction IDs: reports are sent " + "without them and are not retried for the rest of this client's life. Upgrade " + "Kaiten to make retries safe." + ) + response = self._send_report(path, body(None), keyed=False) + + return UsageReportResult( + usage=self._api.validate(EntitlementUsage, self._api.decode(response), response), + replayed=response.headers.get("idempotent-replayed") == "true", + metadata_dropped=bool(response.headers.get("kaiten-metadata-dropped")), ) - return self._api.send( - "POST", - api_path("instances", instance_slug, "entitlements", entitlement_slug, "usage"), - EntitlementUsage, - body=body, - errors={409: ThresholdExceededError}, + + def _send_report(self, path: str, body: dict[str, Any], *, keyed: bool) -> Any: + try: + return self._api.request( + "POST", path, body=body, idempotent=keyed, errors={409: ThresholdExceededError} + ) + except ThresholdExceededError as error: + # The usage endpoint answers 409 for two reasons; only the code tells them apart. + if error.code == _TRANSACTION_ID_REUSED: + raise TransactionIdReusedError( + str(error), response=error.response, problem=error.problem + ) from None + raise + + def list_usage_reports( + self, + instance_slug: str, + entitlement_slug: str, + *, + from_: datetime | str | None = None, + to: datetime | str | None = None, + transaction_id: str | None = None, + limit: int | None = None, + ) -> builtins.list[UsageReport]: + """Return the usage history of one entitlement for one instance, oldest first. + + Every report the API accepted in the range, with the counter before and after it and + the limit it was gated on. Decimals are strings, exact to the digit. Pages are walked + until the range is exhausted or ``limit`` reports are held. + + Args: + instance_slug: The instance. + entitlement_slug: The entitlement. + from_: Start of the range, inclusive. Defaults to 30 days before ``to``; a value + before the start of the organization's usage history is refused with the code + ``ListUsageReports.OutsideRetention``. + to: End of the range, exclusive. Defaults to now. + transaction_id: Only the report sent with this key. + limit: The most reports to return. + """ + if limit is not None and limit < 1: + raise ValueError(f"limit must be at least 1, got {limit}") + path = api_path( + "instances", instance_slug, "entitlements", entitlement_slug, "usage", "reports" + ) + reports: builtins.list[UsageReport] = [] + after_seq: int | None = None + while True: + wanted = _USAGE_HISTORY_PAGE if limit is None else limit - len(reports) + size = min(_USAGE_HISTORY_PAGE, wanted) + page = self._api.get( + path, + UsageReportPage, + params={ + "from": optional_datetime(from_), + "to": optional_datetime(to), + "transactionId": transaction_id, + "afterSeq": after_seq, + "limit": size, + }, + ) + items = page.items or [] + reports.extend(items) + # The API names a next page only after a full one; a short page is the last. + if ( + page.next_after_seq is None + or len(items) < size + or (limit is not None and len(reports) >= limit) + ): + return reports + if after_seq is not None and page.next_after_seq <= after_seq: + raise PaginationError( + f"GET {path}: nextAfterSeq {page.next_after_seq} does not move past " + f"{after_seq}; refusing to return a partial history as a complete one" + ) + after_seq = page.next_after_seq + + def export_usage_reports( + self, + instance_slug: str, + entitlement_slug: str, + destination: BinaryIO, + *, + from_: datetime | str | None = None, + to: datetime | str | None = None, + format: UsageExportFormat = "csv", + ) -> int: + """Stream the usage history of one entitlement for one instance into ``destination``. + + CSV (with a header row) or NDJSON (one report per line), written as it arrives, so the + export is never held whole in memory. At most 366 days per call. Returns the number of + bytes written. + + Args: + instance_slug: The instance. + entitlement_slug: The entitlement. + destination: A binary file-like object, such as ``open(path, "wb")``. + from_: Start of the range, inclusive (see :meth:`list_usage_reports`). + to: End of the range, exclusive. Defaults to now. + format: ``"csv"`` (the default) or ``"json"`` for NDJSON. + """ + return self._api.stream_to( + api_path( + "instances", + instance_slug, + "entitlements", + entitlement_slug, + "usage", + "reports", + "export", + ), + destination, + params={ + "from": optional_datetime(from_), + "to": optional_datetime(to), + "format": format, + }, + ) + + def export_organization_usage_reports( + self, + destination: BinaryIO, + *, + from_: datetime | str | None = None, + to: datetime | str | None = None, + format: UsageExportFormat = "csv", + instance_slug: str | None = None, + instance_id: str | None = None, + entitlement_slug: str | None = None, + entitlement_id: str | None = None, + ) -> int: + """Stream the whole organization's usage history into ``destination``. + + Every instance and entitlement, at most 31 days per call, optionally for one instance + or one entitlement. ``instance_id`` and ``entitlement_id`` reach deleted ones, whose + reports are kept. Returns the number of bytes written. + + Args: + destination: A binary file-like object, such as ``open(path, "wb")``. + from_: Start of the range, inclusive (see :meth:`list_usage_reports`). + to: End of the range, exclusive. Defaults to now. + format: ``"csv"`` (the default) or ``"json"`` for NDJSON. + instance_slug: Only this instance. + instance_id: Only this instance, by id. + entitlement_slug: Only this entitlement. + entitlement_id: Only this entitlement, by id. + """ + return self._api.stream_to( + api_path("usage", "reports", "export"), + destination, + params={ + "from": optional_datetime(from_), + "to": optional_datetime(to), + "format": format, + "instanceSlug": instance_slug, + "instanceId": instance_id, + "entitlementSlug": entitlement_slug, + "entitlementId": entitlement_id, + }, ) def get_integration(self, instance_slug: str, integration_name: str) -> InstanceIntegration: @@ -347,3 +589,14 @@ def _instance( } ), } + + +def _refuses_transaction_id(error: APIStatusError) -> bool: + """Whether ``error`` is an API without transaction IDs refusing the field. + + That is a validation problem located on the field, from an API that does not know the code + newer ones answer a malformed key with. The key was checked before sending. + """ + if error.status_code not in (400, 422) or error.code == _INVALID_TRANSACTION_ID: + return False + return any(detail.location == "body.transactionId" for detail in error.errors) diff --git a/src/kaitencloud/types/__init__.py b/src/kaitencloud/types/__init__.py index c128fc0..2da5562 100644 --- a/src/kaitencloud/types/__init__.py +++ b/src/kaitencloud/types/__init__.py @@ -83,6 +83,8 @@ Targetings, Token, UsageBehavior, + UsageReport, + UsageReportPage, UsageReportStatus, User, Variant, @@ -190,6 +192,8 @@ "Targetings", "Token", "UsageBehavior", + "UsageReport", + "UsageReportPage", "UsageReportStatus", "User", "Variant", diff --git a/src/kaitencloud/types/_core.py b/src/kaitencloud/types/_core.py index c09dbd4..356eae7 100644 --- a/src/kaitencloud/types/_core.py +++ b/src/kaitencloud/types/_core.py @@ -86,6 +86,8 @@ "Targetings", "Token", "UsageBehavior", + "UsageReport", + "UsageReportPage", "UsageReportStatus", "User", "Variant", @@ -1478,6 +1480,81 @@ class Token(KaitenModel): """ +class UsageReport(KaitenModel): + aggregation_method: str = Field(alias="aggregationMethod") + """The entitlement's aggregation method when the report was accepted""" + + behavior: UsageBehavior | str + """append adds the value through the aggregation method; set overwrites the counter""" + + delta: str + """valueAfter minus valueBefore; negative for a set that lowered the counter""" + + entitlement_id: str = Field(alias="entitlementId") + """The entitlement the report was made for. It may since have been deleted.""" + + event_count_after: int = Field(alias="eventCountAfter") + """Reports counted in the window after this one""" + + instance_id: str = Field(alias="instanceId") + """The instance the report was made for. It may since have been deleted.""" + + license_id: str = Field(alias="licenseId") + """The instance's licence when the report was accepted""" + + limit_value: str | None = Field(default=None, alias="limitValue") + """The limit in force when the report was accepted. Null when unlimited.""" + + overage_delta: str = Field(alias="overageDelta") + """How much the usage above limitValue moved: max(0, valueAfter - limitValue) - max(0, + valueBefore - limitValue). 0 when unlimited. + """ + + overage_percent: int | None = Field(default=None, alias="overagePercent") + """The overage allowed above limitValue, in percent, when the report was accepted. Null + when unlimited. + """ + + properties: dict[str, Any] | None = None + """The report's metadata, when it was stored""" + + report_seq: int = Field(alias="reportSeq") + """Position of the report among the pair's accepted reports, from 1, without gaps""" + + reported_at: datetime = Field(alias="reportedAt") + """When the server accepted the report (UTC, millisecond precision)""" + + reported_value: str = Field(alias="reportedValue") + """The value as sent: a delta for append, an absolute value for set""" + + transaction_id: str | None = Field(default=None, alias="transactionId") + """The report's idempotency key, when it was sent with one""" + + value_after: str = Field(alias="valueAfter") + """The counter after the report""" + + value_before: str = Field(alias="valueBefore") + """The counter before the report, after any window reset""" + + window_end: datetime | None = Field(default=None, alias="windowEnd") + """End of the usage window the report counted in (exclusive). Null for a lifetime + entitlement. + """ + + window_start: datetime | None = Field(default=None, alias="windowStart") + """Start of the usage window the report counted in (inclusive). Null for a lifetime + entitlement. + """ + + +class UsageReportPage(KaitenModel): + items: list[UsageReport] + """The reports, in reportSeq order""" + + next_after_seq: int | None = Field(default=None, alias="nextAfterSeq") + """Pass as afterSeq to read the next page. Absent on the last page.""" + + class User(KaitenModel): """A reference to the user behind an action.""" @@ -1575,6 +1652,8 @@ class Variant(KaitenModel): ServiceAccount, SystemTokenIssuance, Token, + UsageReport, + UsageReportPage, User, Variant, ): diff --git a/tests/test_contract.py b/tests/test_contract.py index 0ce28e0..d6002aa 100644 --- a/tests/test_contract.py +++ b/tests/test_contract.py @@ -25,6 +25,7 @@ from __future__ import annotations import asyncio +import io from collections.abc import Callable from dataclasses import dataclass from datetime import datetime, timedelta, timezone @@ -386,6 +387,42 @@ class Case: ), minimal=lambda c: c.instances.report_usage("acme-production", "seats", 1), ), + Case( + "reportEntitlementUsageMetric", + lambda c: c.instances.report_usage_detailed( + "acme-production", "tokens", 1200, transaction_id="llm-call:9f2c:tokens" + ), + required=frozenset({"transactionId"}), + name="reportEntitlementUsageMetric-keyed", + ), + Case( + "listUsageReports", + lambda c: c.instances.list_usage_reports( + "acme-production", "tokens", from_=NOW, to=LATER, transaction_id="evt-1", limit=10 + ), + minimal=lambda c: c.instances.list_usage_reports("acme-production", "tokens"), + ), + Case( + "exportUsageReports", + lambda c: c.instances.export_usage_reports( + "acme-production", "tokens", io.BytesIO(), from_=NOW, to=LATER, format="json" + ), + minimal=lambda c: c.instances.export_usage_reports( + "acme-production", "tokens", io.BytesIO() + ), + ), + Case( + "exportOrganizationUsageReports", + lambda c: c.instances.export_organization_usage_reports( + io.BytesIO(), + from_=NOW, + to=LATER, + format="csv", + instance_slug="acme-production", + entitlement_id="22222222-2222-2222-2222-222222222222", + ), + minimal=lambda c: c.instances.export_organization_usage_reports(io.BytesIO()), + ), Case( "get-instance-integration", lambda c: c.instances.get_integration("acme-production", "attio"), diff --git a/tests/test_usage_reports.py b/tests/test_usage_reports.py new file mode 100644 index 0000000..bc6bc3e --- /dev/null +++ b/tests/test_usage_reports.py @@ -0,0 +1,277 @@ +# Copyright 2026 KAITEN INC +# SPDX-License-Identifier: Apache-2.0 + +"""Usage report keys, and the usage history. + +A ``transaction_id`` makes a report idempotent on the server, so it is the one thing that +makes retrying a report safe: a report is retried with a key and never without one. A key +already used for another report is a client bug, refused with its own error. +""" + +from __future__ import annotations + +import asyncio +import io +import logging +from collections.abc import Callable +from typing import Any + +import httpx +import pytest +from contract import sample +from helpers import Recorder, body, respond + +from kaitencloud import ( + APIStatusError, + PaginationError, + ServiceUnavailableError, + ThresholdExceededError, + TransactionIdReusedError, + UnprocessableEntityError, +) + +USAGE = sample({"$ref": "#/components/schemas/EntitlementUsage"}, "core") +REPORT = sample({"$ref": "#/components/schemas/UsageReport"}, "core") + + +def problem(status: int, code: str, **extra: Any) -> httpx.Response: + return httpx.Response( + status, + json={"title": "Error", "status": status, "code": code, **extra}, + headers={"content-type": "application/problem+json"}, + ) + + +def in_turn(*answers: Callable[[httpx.Request], httpx.Response]) -> Callable[..., httpx.Response]: + """Answer the nth request with the nth answer, the last one repeating.""" + calls = {"n": 0} + + def handler(request: httpx.Request) -> httpx.Response: + answer = answers[min(calls["n"], len(answers) - 1)] + calls["n"] += 1 + return answer(request) + + return handler + + +def test_the_key_is_sent_only_when_given(make_client: Callable[..., Any]) -> None: + client, recorder = make_client(respond(200, USAGE)) + client.instances.report_usage("acme", "tokens", 1200, transaction_id="llm-call:9f2c:tokens") + client.instances.report_usage("acme", "tokens", 1) + assert body(recorder[0])["transactionId"] == "llm-call:9f2c:tokens" + assert "transactionId" not in body(recorder[1]) + + +def test_a_keyed_report_is_retried_under_the_same_key( + make_client: Callable[..., Any], delays: list[float] +) -> None: + client, recorder = make_client( + in_turn( + lambda _: problem(503, "ReportEntitlementUsageMetric.LedgerUnavailable"), + lambda _: httpx.Response(200, json=USAGE), + ) + ) + client.instances.report_usage("acme", "tokens", 1200, transaction_id="evt-1") + assert [body(request)["transactionId"] for request in recorder] == ["evt-1", "evt-1"] + assert len(delays) == 1 + + +def test_an_unkeyed_report_is_never_retried(make_client: Callable[..., Any]) -> None: + client, recorder = make_client(respond(503, {"title": "Unavailable"})) + with pytest.raises(ServiceUnavailableError): + client.instances.report_usage("acme", "tokens", 1) + assert len(recorder) == 1 + + +@pytest.mark.parametrize( + ("code", "expected", "not_expected"), + [ + ( + "ReportEntitlementUsageMetric.TransactionIdReused", + TransactionIdReusedError, + ThresholdExceededError, + ), + ( + "ReportEntitlementUsageMetric.ThresholdExceeded", + ThresholdExceededError, + TransactionIdReusedError, + ), + ], +) +def test_a_conflict_is_named_by_its_code( + make_client: Callable[..., Any], + code: str, + expected: type[APIStatusError], + not_expected: type[APIStatusError], +) -> None: + client, recorder = make_client( + lambda _: problem( + 409, code, errors=[{"location": "body.transactionId", "value": {"reportSeq": 2}}] + ) + ) + with pytest.raises(expected) as raised: + client.instances.report_usage("acme", "tokens", 1, transaction_id="evt-1") + assert not isinstance(raised.value, not_expected) + assert raised.value.code == code + assert len(recorder) == 1, "a conflict is never retried" + + +@pytest.mark.parametrize("key", ["has space", "x" * 129, "a/b", ""]) +def test_a_malformed_key_is_refused_before_sending( + make_client: Callable[..., Any], key: str +) -> None: + client, recorder = make_client(respond(200, USAGE)) + with pytest.raises(ValueError, match="transaction_id"): + client.instances.report_usage("acme", "tokens", 1, transaction_id=key) + assert len(recorder) == 0 + + +def test_the_detailed_report_reads_what_the_server_said(make_client: Callable[..., Any]) -> None: + client, _ = make_client( + respond( + 200, + USAGE, + headers={"Idempotent-Replayed": "true", "Kaiten-Metadata-Dropped": "too_large"}, + ) + ) + result = client.instances.report_usage_detailed("acme", "tokens", 1, transaction_id="evt-1") + assert result.replayed + assert result.metadata_dropped + + plain, _ = make_client(respond(200, USAGE)) + result = plain.instances.report_usage_detailed("acme", "tokens", 1) + assert not result.replayed + assert not result.metadata_dropped + + +def test_an_api_without_keys_gets_the_report_resent_without_one( + make_client: Callable[..., Any], caplog: pytest.LogCaptureFixture +) -> None: + refused = { + "title": "Unprocessable Entity", + "status": 422, + "code": "UNPROCESSABLE", + "errors": [{"location": "body.transactionId", "message": "unexpected property"}], + } + + def handler(request: httpx.Request) -> httpx.Response: + if "transactionId" in (body(request) or {}): + return httpx.Response(422, json=refused) + return httpx.Response(200, json=USAGE) + + client, recorder = make_client(handler) + with caplog.at_level(logging.WARNING, logger="kaitencloud"): + client.instances.report_usage("acme", "tokens", 1, transaction_id="evt-1") + client.instances.report_usage("acme", "tokens", 2, transaction_id="evt-2") + + assert len(recorder) == 3 + assert "transactionId" not in body(recorder[1]) + assert "transactionId" not in body(recorder[2]), "later reports go out without a key" + warnings = [record for record in caplog.records if "transaction IDs" in record.getMessage()] + assert len(warnings) == 1 + + +def test_a_new_apis_malformed_key_code_is_not_an_old_api(make_client: Callable[..., Any]) -> None: + client, recorder = make_client( + lambda _: problem( + 422, + "ReportEntitlementUsageMetric.InvalidTransactionId", + errors=[{"location": "body.transactionId", "message": "invalid"}], + ) + ) + with pytest.raises(UnprocessableEntityError): + client.instances.report_usage("acme", "tokens", 1, transaction_id="evt-1") + assert len(recorder) == 1 + + +def test_the_async_client_retries_a_keyed_report_too(make_async_client: Callable[..., Any]) -> None: + async def scenario() -> Recorder: + client, recorder = make_async_client( + in_turn(lambda _: httpx.Response(502), lambda _: httpx.Response(200, json=USAGE)) + ) + async with client: + result = await client.instances.report_usage_detailed( + "acme", "tokens", 1, transaction_id="evt-1" + ) + assert not result.replayed + return recorder + + assert len(asyncio.run(scenario())) == 2 + + +def page(seqs: list[int], next_after_seq: int | None) -> dict[str, Any]: + payload: dict[str, Any] = {"items": [{**REPORT, "reportSeq": seq} for seq in seqs]} + if next_after_seq is not None: + payload["nextAfterSeq"] = next_after_seq + return payload + + +def test_the_history_walks_its_pages( + make_client: Callable[..., Any], monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setattr("kaitencloud._sync.resources.instances._USAGE_HISTORY_PAGE", 2) + client, recorder = make_client( + in_turn( + lambda _: httpx.Response(200, json=page([1, 2], 2)), + lambda _: httpx.Response(200, json=page([3], None)), + ) + ) + reports = client.instances.list_usage_reports("acme", "tokens", from_="2026-10-01T00:00:00Z") + assert [report.report_seq for report in reports] == [1, 2, 3] + assert recorder[0].url.path.endswith("/instances/acme/entitlements/tokens/usage/reports") + assert "afterSeq" not in recorder[0].url.params + assert recorder[1].url.params["afterSeq"] == "2" + assert recorder[0].url.params["from"] == "2026-10-01T00:00:00Z" + + +def test_the_history_stops_at_its_limit(make_client: Callable[..., Any]) -> None: + client, recorder = make_client(respond(200, page([1, 2], 2))) + reports = client.instances.list_usage_reports("acme", "tokens", limit=2) + assert len(reports) == 2 + assert recorder.only.url.params["limit"] == "2" + + +def test_a_cursor_that_does_not_move_is_refused( + make_client: Callable[..., Any], monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setattr("kaitencloud._sync.resources.instances._USAGE_HISTORY_PAGE", 2) + client, _ = make_client(respond(200, page([1, 2], 2))) + with pytest.raises(PaginationError): + client.instances.list_usage_reports("acme", "tokens") + + +def test_an_export_streams_into_its_destination(make_client: Callable[..., Any]) -> None: + csv = b"organization_id,instance_id\no,i\n" + client, recorder = make_client(respond(200, csv, headers={"content-type": "text/csv"})) + destination = io.BytesIO() + written = client.instances.export_usage_reports("acme", "tokens", destination, format="json") + assert written == len(csv) + assert destination.getvalue() == csv + assert recorder.only.url.params["format"] == "json" + + +def test_an_export_refusal_raises_with_its_code(make_client: Callable[..., Any]) -> None: + client, _ = make_client(lambda _: problem(422, "ExportUsageReports.OutsideRetention")) + destination = io.BytesIO() + with pytest.raises(UnprocessableEntityError) as raised: + client.instances.export_organization_usage_reports( + destination, from_="2025-01-01T00:00:00Z" + ) + assert raised.value.code == "ExportUsageReports.OutsideRetention" + assert destination.getvalue() == b"" + + +def test_an_export_retries_an_unavailable_api_before_writing( + make_async_client: Callable[..., Any], +) -> None: + async def scenario() -> tuple[int, Recorder, bytes]: + client, recorder = make_async_client( + in_turn(lambda _: httpx.Response(503), lambda _: httpx.Response(200, content=b"a,b\n")) + ) + destination = io.BytesIO() + async with client: + written = await client.instances.export_organization_usage_reports(destination) + return written, recorder, destination.getvalue() + + written, recorder, content = asyncio.run(scenario()) + assert (written, len(recorder), content) == (4, 2, b"a,b\n")