diff --git a/.github/scripts/check_public_api.py b/.github/scripts/check_public_api.py index 702abf170..bc8a3a8d6 100644 --- a/.github/scripts/check_public_api.py +++ b/.github/scripts/check_public_api.py @@ -59,6 +59,15 @@ def _is_public_alias(obj: object) -> bool: return _is_public_posthog_path(target_path) +def _is_private_sdk_alias(obj: object) -> bool: + target = _alias_target_path(obj) + return ( + target.startswith("posthog.") + and not _is_excluded_path(target) + and not _is_public_posthog_path(target) + ) + + def _module_exports(module: object) -> set[str] | None: exports = getattr(module, "exports", None) if exports is None: @@ -98,7 +107,9 @@ def _is_public_member(parent: object, name: str, member: object) -> bool: return False if _is_alias(member): - return _is_public_alias(member) + return _is_public_alias(member) or ( + exports is not None and _is_private_sdk_alias(member) + ) if _is_module(member): return _is_public_posthog_path(_object_path(member)) @@ -119,7 +130,7 @@ def _signature(obj: object) -> str | None: def _signature_for_path(obj: object) -> str | None: signature = _signature(obj) - if signature is None: + if signature is None or (not signature and _is_alias(obj)): return None name = str(getattr(obj, "name", "")) @@ -144,7 +155,7 @@ def _attribute_details(obj: object) -> str: def _record(obj: object) -> str: - if _is_alias(obj): + if _is_alias(obj) and not _is_private_sdk_alias(obj): return f"alias {_object_path(obj)} -> {_alias_target_path(obj)}" if _is_module(obj): @@ -169,15 +180,21 @@ def _record(obj: object) -> str: return f"{getattr(obj, 'kind', 'object')} {_object_path(obj)}" -def _iter_class_members(cls: object) -> Iterable[object]: +def _iter_class_members( + cls: object, *, exported_alias: bool = False +) -> Iterable[object]: members = getattr(cls, "members", {}) for name, member in sorted(members.items()): - if not _is_public_member(cls, str(name), member): + if not ( + _is_public_name(str(name)) + if exported_alias + else _is_public_member(cls, str(name), member) + ): continue yield member - if not _is_alias(member) and _is_class(member): - yield from _iter_class_members(member) + if _is_class(member) and (exported_alias or not _is_alias(member)): + yield from _iter_class_members(member, exported_alias=exported_alias) def _iter_module_members(module: object) -> Iterable[object]: @@ -188,6 +205,8 @@ def _iter_module_members(module: object) -> Iterable[object]: if _is_alias(member): if _is_public_member(module, str(name), member): yield member + if _is_private_sdk_alias(member) and _is_class(member): + yield from _iter_class_members(member, exported_alias=True) continue if _is_module(member): diff --git a/.github/scripts/check_strict_types.sh b/.github/scripts/check_strict_types.sh index bc95aa8d7..7b46c0af0 100755 --- a/.github/scripts/check_strict_types.sh +++ b/.github/scripts/check_strict_types.sh @@ -15,6 +15,14 @@ import atexit import posthog from posthog import FeatureFlagEvaluations, FlagValue, Posthog +from posthog.ai.evaluations import ( + AsyncOfflineEvaluations, + BooleanScorerConfig, + EvaluationItem, + NumericPassingRule, + NumericScorerConfig, + OfflineEvaluations, +) client = Posthog("phc_test") atexit.register(client.shutdown) @@ -31,6 +39,40 @@ active: posthog.Span | None = posthog.get_active_span() span.end() _ = (flag_value, all_flags, enabled, payload, evaluations, active) + + +async def offline_evaluation_types( + sync_client: OfflineEvaluations, async_client: AsyncOfflineEvaluations +) -> None: + numeric = sync_client.scorers.create( + name="Relevance", + kind="numeric", + config=NumericScorerConfig( + min=0, max=1, + passing_rule=NumericPassingRule(operator="gte", threshold=0.8), + ), + ) + config: NumericScorerConfig = numeric.config + experiment = sync_client.create_experiment(name="baseline", run_source="ci") + item = EvaluationItem(input="question", output="answer") + experiment.upload_result(item=item, scorer_version_id=numeric.current_version_id, value=0.9) + boolean = await async_client.scorers.create( + name="Toxicity", kind="boolean", config=BooleanScorerConfig(true_is_failure=True) + ) + async_experiment = await async_client.create_experiment(name="async run") + await async_experiment.upload_result(item=item, scorer_version_id=boolean.current_version_id, value=False) + _ = config +PY + +cat > "$tmp/invalid_offline_types.py" <<'PY' +# pyright: strict +from posthog.ai.evaluations import BooleanScorerConfig, CategoricalScorerConfig, NumericPassingRule, NumericScorerConfig, OfflineEvaluations + +client = OfflineEvaluations(project_id=123, secret_key="phx_example") +NumericPassingRule(operator="gt", threshold=0.8) # expected-type-error +NumericScorerConfig(unknown_setting=True) # expected-type-error +CategoricalScorerConfig(selection_mode="single") # expected-type-error +client.scorers.create(name="Bad config", kind="numeric", config=BooleanScorerConfig(true_is_failure=True)) # expected-type-error PY "$tmp/.venv/bin/python" - <<'PY' > "$tmp/public_api_access.py" @@ -73,3 +115,21 @@ JSON cd "$tmp" "$tmp/.venv/bin/python" -m pyright strict_posthog_types.py public_api_access.py + +if "$tmp/.venv/bin/python" -m pyright --outputjson invalid_offline_types.py > invalid_types.json; then + echo "Invalid offline evaluation configurations unexpectedly passed type checking." + exit 1 +fi +"$tmp/.venv/bin/python" - <<'PY' +import json +from pathlib import Path + +expected = { + index for index, line in enumerate(Path("invalid_offline_types.py").read_text().splitlines()) + if "expected-type-error" in line +} +diagnostics = json.loads(Path("invalid_types.json").read_text())["generalDiagnostics"] +actual = {item["range"]["start"]["line"] for item in diagnostics if item["severity"] == "error"} +assert actual == expected, diagnostics +print("Offline evaluation configuration types reject invalid inputs.") +PY diff --git a/.github/scripts/test_check_public_api.py b/.github/scripts/test_check_public_api.py index c1bf97f25..5aa6da473 100644 --- a/.github/scripts/test_check_public_api.py +++ b/.github/scripts/test_check_public_api.py @@ -6,6 +6,7 @@ import importlib.util from pathlib import Path from types import SimpleNamespace +from tempfile import TemporaryDirectory SCRIPT_PATH = Path(__file__).with_name("check_public_api.py") @@ -42,8 +43,41 @@ def test_attribute_details_uses_placeholder_values() -> None: ) +def test_explicit_exports_from_private_modules_include_signatures() -> None: + import griffe + + check_public_api = load_check_public_api() + with TemporaryDirectory() as directory: + package = Path(directory) / "posthog" + package.mkdir() + (package / "__init__.py").write_text( + 'from ._implementation import Client, Hidden\n__all__ = ["Client"]\n' + ) + (package / "_implementation.py").write_text( + "class Client:\n" + " def __init__(self, *, key: str): pass\n" + " def upload(self, value: float) -> bool: return True\n" + " def _internal(self): pass\n" + "class Hidden: pass\n" + ) + module = griffe.load( + "posthog", + search_paths=[directory], + allow_inspection=False, + try_relative_path=False, + ) + records = [ + check_public_api._record(obj) + for obj in check_public_api._iter_module_members(module) + ] + assert "class posthog.Client(*, key: str)" in records + assert "method posthog.Client.upload(value: float)" in records + assert not any("Hidden" in record or "_internal" in record for record in records) + + def main() -> int: test_attribute_details_uses_placeholder_values() + test_explicit_exports_from_private_modules_include_signatures() print("check_public_api tests passed.") return 0 diff --git a/.sampo/changesets/offline-evaluations.md b/.sampo/changesets/offline-evaluations.md new file mode 100644 index 000000000..b47fe83ef --- /dev/null +++ b/.sampo/changesets/offline-evaluations.md @@ -0,0 +1,5 @@ +--- +pypi/posthog: minor +--- + +Add typed synchronous and asynchronous offline evaluation clients with reusable items, acknowledged bulk uploads, resumable identities, and explicit experiment completion. Manage scorers and immutable versions through the SDK, including boolean, numeric, and categorical polarity configuration. diff --git a/examples/example-offline-evaluations/README.md b/examples/example-offline-evaluations/README.md new file mode 100644 index 000000000..c1ec2207f --- /dev/null +++ b/examples/example-offline-evaluations/README.md @@ -0,0 +1,196 @@ +# Offline evaluations + +Upload evaluations you have already run, manage reusable scorers, and explicitly finish an experiment. The SDK returns server acknowledgments for every operation. It does not execute your application or judges, create traces, or upload in the background. + +## Credentials and setup + +Offline evaluations must be enabled for your PostHog project. Use the project's numeric ID and the app host (`https://us.posthog.com`, `https://eu.posthog.com`, or your self-hosted URL). + +| Operation | Credential and scopes | +| --- | --- | +| Create experiments, upload results, complete/fail | Personal API key or project secret key with `offline_evaluation_ingestion:write` | +| Fetch scorers and versions | Personal API key with `llm_analytics:read` | +| Create/update scorers and create versions | Personal API key with `llm_analytics:write` | + +A personal key can cover both scorer management and uploads. Project secret keys can upload using a known scorer-version UUID without scorer read access. All operations use the explicit project/environment; scorer names are not unique. + +From this repository's root: + +```bash +export POSTHOG_PROJECT_ID=123 +export POSTHOG_SECRET_KEY='your-secret-key' +export POSTHOG_HOST='https://us.posthog.com' + +# Optional one-time provisioning; requires a personal key. +uv run python examples/example-offline-evaluations/main.py --create-scorer + +# Use the stable scorer ID printed above; its current version is fetched once. +export POSTHOG_SCORER_ID='your-scorer-uuid' +uv run python examples/example-offline-evaluations/main.py --state capital-run.json +``` + +Alternatively, set `POSTHOG_SCORER_VERSION_ID` to an existing **boolean** version UUID. This takes precedence over `POSTHOG_SCORER_ID` and avoids discovery. The synthetic example stores its experiment, items, and results in the state file before uploading. Rerunning with that file replays the same experiment; a new file starts a new experiment. The state includes evaluation data and never includes your API key. No LLM provider or API key is needed for this example. + +## Upload a run + +```python +from posthog.ai.evaluations import ( + EvaluationItem, + EvaluationResult, + OfflineEvaluations, +) + +with OfflineEvaluations( + project_id=123, + secret_key=secret_key, + host="https://us.posthog.com", +) as evaluations: + scorer = evaluations.scorers.get(scorer_id) + version_id = scorer.current_version_id # Pin once for this run. + experiment = evaluations.create_experiment( + name="support-agent", + run_source="ci", # Optional: ci, local, or scheduled. + expected_item_count=1, + expected_result_count=1, + application_version="your-commit-sha", + ) + item = EvaluationItem( + input="What is the capital of France?", + output="Paris", + expected_output="Paris", + case_key="capital-france", + ) + receipt = experiment.upload_result( + item=item, + scorer_version_id=version_id, + value=True, + ) + print(receipt.results[0].id, receipt.results[0].created) + experiment.complete() +``` + +Use one `EvaluationItem` for the application's execution and reuse it across scorer results. `upload_results([EvaluationResult(...), ...])` uploads a finite sequence, eagerly and sequentially. It splits at 1,000 items/results or 5 MiB of UTF-8 JSON per request, whichever comes first. Item payloads are limited to 1 MiB; result payloads to 256 KiB. Payload JSON supports at most 32 nested containers. The SDK raises for invalid values and oversized payloads without truncating them. Memory use scales with the sequence you supply, so submit caller-sized batches for large datasets. + +`status="ok"` means the evaluator ran successfully. Its boolean, number, or category-list `value` is uploaded with its original polarity. Numeric scores use finite double-precision (binary64) values; integer inputs are rejected if conversion would change their value. Use `error`, `skipped`, or `not_applicable` without a value for other outcomes. `error_code` and a non-null `error_message` are valid only with `status="error"`. + +Omitted payload fields, explicit `None`, and empty payloads are distinct: + +```python +EvaluationItem() # No payload declaration. +EvaluationItem(payload={}) # Explicitly empty payload. +EvaluationItem(input=None) # Payload contains a null input. +``` + +The client context manager releases connections. Call `complete()` explicitly after uploads succeed, or `fail()` when you deliberately want to close the experiment as failed. Neither exceptions nor async cancellation automatically close an experiment. Expected counts include every outcome status; the server checks them on completion. + +## Typed scorer configuration and polarity + +Configurations are public `TypedDict` types: editors suggest fields and type checkers connect `create(kind=...)` to the correct configuration. They are ordinary dictionaries at runtime; server validation remains authoritative. + +```python +from posthog.ai.evaluations import ( + BooleanScorerConfig, + CategoricalPassingRule, + CategoricalScorerConfig, + CategoricalScorerOption, + NumericPassingRule, + NumericScorerConfig, +) + +boolean = BooleanScorerConfig( + true_is_failure=True, + true_label="Unsafe", + false_label="Safe", +) +numeric = NumericScorerConfig( + min=0, + max=1, + step=0.01, + passing_rule=NumericPassingRule(operator="gte", threshold=0.8), +) +categorical = CategoricalScorerConfig( + options=[ + CategoricalScorerOption(key="safe", label="Safe"), + CategoricalScorerOption(key="helpful", label="Helpful"), + ], + selection_mode="multiple", + min_selections=1, + max_selections=2, + passing_rule=CategoricalPassingRule(categories=["safe", "helpful"]), +) +scorer = evaluations.scorers.create( + name="Answer relevance", + kind="numeric", + config=numeric, +) +version_id = scorer.current_version_id +``` + +Boolean `true_is_failure=True` makes false pass; false, null, or omission makes true pass. Numeric passing rules support `gte` and `lte`; the finite threshold must be within any configured bounds. Numeric bounds, step, and passing rule can be omitted or null. For categorical scores, **every** returned category must be in the passing list. Selection defaults to `single`; minimum and maximum selections apply only to `multiple`. Numeric/categorical scores are neutral when their passing rule is omitted or null. Polarity is pinned to each immutable version. + +Manage metadata independently from versioned configuration: + +```python +scorer = evaluations.scorers.get(scorer_id) +evaluations.scorers.update(scorer.id, description="Updated description", archived=False) +updated = evaluations.scorers.create_version( + scorer.id, + config=scorer.config, + base_version=scorer.current_version, +) +version = evaluations.scorers.get_version(scorer.id, updated.current_version_id) +``` + +`create_version()` returns the updated scorer. Identical configuration is allowed when, for example, your judging prompt changes. Supply `base_version` for optimistic concurrency; a stale version raises a 409 error without silently refreshing and retrying. Omit metadata fields to preserve them; `description=None` clears the description. Scorer kind cannot change. + +`scorers.list(limit=100, offset=0, search=..., kind=..., archived=..., order_by=...)` returns a page with `results`, `count`, and `next_offset`; preserve the filters when continuing. Omitted `archived` lists active scorers. `list_versions(scorer_id, limit=50, cursor=...)` returns newest versions first and a `next_cursor`; its maximum limit is 100. Stop when the corresponding continuation is `None`. + +Scorer creation, version creation, and metadata updates are never automatically retried. If a response is lost, inspect server state before repeating a mutation. Version creation is not made idempotent by `base_version`. + +## Receipts, errors, and recovery + +Upload receipts contain persisted item/result UUIDs, original acceptance times, and `created` flags. `created=False` acknowledges an exact duplicate. A bulk upload returns `chunks`; each chunk commits independently. A malformed or incomplete acknowledgment raises rather than reporting success. + +`EvaluationAPIError` exposes `status`, `code`, `detail`, `attr`, `errors`, `response`, `retry_after`, and `persistence`: + +- `rejected`: the server rejected the request. +- `unknown`: it may have committed, for example if its response was lost. +- `not_sent`: it was not sent. + +Errors do not include evaluation content or credentials in their message; inspect their structured attributes explicitly when needed. `BulkUploadError` also exposes `completed.chunks`, `failed_index`, the exact `failed_chunk`, and `pending_results`. After resolving the cause, retry the pending results against the same experiment. An uncertain chunk can safely be replayed with identical identities and content. Do not generate new item IDs during recovery. + +```python +from posthog.ai.evaluations import BulkUploadError + +try: + receipts = experiment.upload_results(results) +except BulkUploadError as error: + acknowledged = error.completed + pending = error.pending_results + # Preserve these and inspect the cause before retrying: + # experiment.upload_results(pending) +``` + +Creation generates an experiment UUID and start time once. A successful handle exposes its full `submission`; creation errors expose `error.submission`. Retry `evaluations.create_experiment(**submission)` with all original fields to recover a failed acknowledgment. `resume_experiment(id)` creates a local handle without a read request; it does not create missing experiments. + +For recovery after process termination, persist explicit experiment `id`, `started_at`, and all creation fields **before** the creation call. Also save every `item.to_dict()` and `result.to_dict()` before upload. Restore items with `EvaluationItem.from_dict(...)`; attach the matching restored item with `EvaluationResult.from_dict(saved_result, item=item)`. This preserves payloads and re-declares items safely if their prior acknowledgment was lost. The example's state file implements this pattern; the SDK itself does not write to disk. + +Default request timeout is 15 seconds with at most three retries. Idempotent offline operations preserve request bytes across retries and honor rate-limit delays. Authentication, validation, and conflicts surface immediately. There is no automatic completion after a failed upload. + +## Async + +Install `posthog[async]` for native async support. Methods and receipts match the synchronous client: + +```python +from posthog.ai.evaluations import AsyncOfflineEvaluations + +async with AsyncOfflineEvaluations(project_id=123, secret_key=secret_key) as evaluations: + scorer = await evaluations.scorers.get(scorer_id) + experiment = await evaluations.create_experiment(name="async-run") + await experiment.upload_result( + item=item, scorer_version_id=scorer.current_version_id, value=True, + ) + await experiment.complete() +``` + +`resume_experiment(id)` is local in both clients and does not require `await`. Use each async client on one event loop. Cancellation propagates; replay the persisted submission when an in-flight upload's outcome is uncertain. diff --git a/examples/example-offline-evaluations/main.py b/examples/example-offline-evaluations/main.py new file mode 100644 index 000000000..a3c34171d --- /dev/null +++ b/examples/example-offline-evaluations/main.py @@ -0,0 +1,171 @@ +"""Run a synthetic offline evaluation, with explicit scorer setup and durable replay.""" + +import argparse +import json +import os +from datetime import datetime, timezone +from pathlib import Path +from typing import Any +from uuid import UUID, uuid4 + +from posthog.ai.evaluations import ( + BooleanScorerConfig, + BulkUploadError, + EvaluationAPIError, + EvaluationItem, + EvaluationResult, + OfflineEvaluations, +) + + +def resolve_version(evaluations: OfflineEvaluations) -> str: + version_id = os.environ.get("POSTHOG_SCORER_VERSION_ID") + if version_id: + return str(UUID(version_id)) + scorer_id = os.environ.get("POSTHOG_SCORER_ID") + if not scorer_id: + raise ValueError( + "Set POSTHOG_SCORER_ID or POSTHOG_SCORER_VERSION_ID, or run --create-scorer once." + ) + scorer = evaluations.scorers.get(scorer_id) + if scorer.kind != "boolean": + raise ValueError("This example requires a boolean scorer.") + return scorer.current_version_id + + +def prepare_state( + evaluations: OfflineEvaluations, project_id: int, host: str +) -> dict[str, Any]: + version_id = resolve_version(evaluations) + cases = [ + { + "key": "france", + "input": "Capital of France?", + "output": "Paris", + "expected": "Paris", + }, + { + "key": "italy", + "input": "Capital of Italy?", + "output": "Milan", + "expected": "Rome", + }, + ] + items = [ + EvaluationItem( + case_key=case["key"], + input=case["input"], + output=case["output"], + expected_output=case["expected"], + ) + for case in cases + ] + results = [ + EvaluationResult( + item=item, + scorer_version_id=version_id, + value=case["output"] == case["expected"], + reasoning="Exact string comparison", + ) + for item, case in zip(items, cases) + ] + return { + "project_id": project_id, + "host": host, + "experiment": { + "id": str(uuid4()), + "name": "Capital city example", + "started_at": datetime.now(timezone.utc).isoformat(), + "run_source": "local", + "expected_item_count": len(items), + "expected_result_count": len(results), + }, + "items": [item.to_dict() for item in items], + "results": [result.to_dict() for result in results], + } + + +def save_state(path: Path, state: dict[str, Any]) -> None: + # Persist identities and content before the first experiment write. Never + # overwrite an existing run or record credentials in its state file. + descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + with os.fdopen(descriptor, "w", encoding="utf-8") as file: + json.dump(state, file, ensure_ascii=False, allow_nan=False, indent=2) + file.flush() + os.fsync(file.fileno()) + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument( + "--create-scorer", + action="store_true", + help="Create a boolean scorer once and print its IDs, then exit.", + ) + parser.add_argument( + "--state", + type=Path, + default=Path("offline-evaluation-state.json"), + help="Create a new prepared run here, or replay the existing run exactly.", + ) + args = parser.parse_args() + project_id = int(os.environ["POSTHOG_PROJECT_ID"]) + host = os.environ.get("POSTHOG_HOST", "https://us.posthog.com") + with OfflineEvaluations( + project_id=project_id, secret_key=os.environ["POSTHOG_SECRET_KEY"], host=host + ) as evaluations: + if args.create_scorer: + scorer = evaluations.scorers.create( + name="Capital exact match", + kind="boolean", + config=BooleanScorerConfig( + true_is_failure=False, true_label="Correct", false_label="Incorrect" + ), + ) + print(f"Scorer ID: {scorer.id}") + print(f"Version ID: {scorer.current_version_id}") + return + + if args.state.exists(): + state = json.loads(args.state.read_text(encoding="utf-8")) + if state["project_id"] != project_id or state["host"] != host: + raise ValueError( + "The saved run belongs to a different project or host." + ) + else: + state = prepare_state(evaluations, project_id, host) + save_state(args.state, state) + + items = {item["id"]: EvaluationItem.from_dict(item) for item in state["items"]} + results = [ + EvaluationResult.from_dict(result, item=items[result["item_id"]]) + for result in state["results"] + ] + try: + experiment = evaluations.create_experiment(**state["experiment"]) + receipt = experiment.upload_results(results) + completed = experiment.complete() + except BulkUploadError as error: + print( + f"Accepted {len(error.completed.chunks)} chunks; stopped at result {error.failed_index} ({error.persistence})." + ) + print( + f"Prepared run retained in {args.state}. Inspect the error before replaying the same file." + ) + raise + except EvaluationAPIError: + print( + f"Prepared run retained in {args.state}; the SDK has not automatically closed the experiment." + ) + raise + print(f"Experiment {experiment.id}: {completed.status}") + print( + f"Acknowledged {sum(len(chunk.results) for chunk in receipt.chunks)} results across {len(receipt.chunks)} chunks." + ) + print( + f"Keep {args.state} for exact replay; choose a different --state path for a new experiment." + ) + + +if __name__ == "__main__": + main() diff --git a/posthog/ai/evaluations/__init__.py b/posthog/ai/evaluations/__init__.py new file mode 100644 index 000000000..7d3e89b0d --- /dev/null +++ b/posthog/ai/evaluations/__init__.py @@ -0,0 +1,88 @@ +"""Typed offline evaluation uploads, experiment lifecycle, and scorer management. + +from posthog.ai.evaluations import OfflineEvaluations, EvaluationItem + +with OfflineEvaluations(project_id=123, secret_key="phs_example") as evaluations: + experiment = evaluations.create_experiment(name="baseline") + item = EvaluationItem(input="Question", output="Answer") + # Use an existing immutable scorer-version UUID for each result. + # experiment.upload_result(item=item, scorer_version_id=version_id, value=True) + # experiment.complete() +""" + +from ._client import ( + AsyncExperiment, + AsyncOfflineEvaluations, + Experiment, + OfflineEvaluations, +) +from ._errors import EvaluationAPIError +from ._scorer_types import ( + BooleanScorerConfig, + CategoricalPassingRule, + CategoricalScorerConfig, + CategoricalScorerOption, + NumericPassingRule, + NumericScorerConfig, + Scorer, + ScorerConfig, + ScorerKind, + ScorerOrder, + ScorerPage, + ScorerVersion, + ScorerVersionPage, +) +from ._scorers import AsyncScorers, Scorers +from ._serialization import JSONValue +from ._types import ( + BulkUploadError, + BulkUploadReceipt, + EvaluationItem, + EvaluationResult, + ExperimentOptions, + ExperimentReceipt, + ItemReceipt, + ResultOptions, + ResultReceipt, + ResultStatus, + RunSource, + ScoreValue, + UploadReceipt, +) + +__all__ = [ + "AsyncExperiment", + "AsyncOfflineEvaluations", + "AsyncScorers", + "BooleanScorerConfig", + "BulkUploadError", + "BulkUploadReceipt", + "CategoricalPassingRule", + "CategoricalScorerConfig", + "CategoricalScorerOption", + "EvaluationAPIError", + "EvaluationItem", + "EvaluationResult", + "Experiment", + "ExperimentOptions", + "ExperimentReceipt", + "ItemReceipt", + "JSONValue", + "NumericPassingRule", + "NumericScorerConfig", + "OfflineEvaluations", + "ResultOptions", + "ResultReceipt", + "ResultStatus", + "RunSource", + "ScoreValue", + "Scorer", + "ScorerConfig", + "ScorerKind", + "ScorerOrder", + "ScorerPage", + "Scorers", + "ScorerVersion", + "ScorerVersionPage", + "UploadReceipt", +] diff --git a/posthog/ai/evaluations/_client.py b/posthog/ai/evaluations/_client.py new file mode 100644 index 000000000..9e7c75c56 --- /dev/null +++ b/posthog/ai/evaluations/_client.py @@ -0,0 +1,493 @@ +from __future__ import annotations + +import asyncio +import json +from dataclasses import dataclass +from datetime import datetime, timezone +from threading import RLock +from types import TracebackType +from typing import Any, Sequence +from uuid import UUID, uuid4 + +from typing_extensions import Unpack + +from ._errors import EvaluationAPIError +from ._scorers import AsyncScorers, Scorers +from ._serialization import ( + _MAX_ITEMS, + _MAX_REQUEST_BYTES, + _MAX_RESULTS, + _count, + _encode, + _identifier, + _timestamp, + _uuid, +) +from ._transport import AsyncTransport, SyncTransport +from ._types import ( + BulkUploadError, + BulkUploadReceipt, + EvaluationItem, + EvaluationResult, + ExperimentOptions, + ExperimentReceipt, + ResultOptions, + UploadReceipt, + _experiment_receipt, + _upload_receipt, +) + + +def _experiment_body(name: str, options: ExperimentOptions) -> bytes: + allowed = set(ExperimentOptions.__annotations__) + if set(options) - allowed: + raise TypeError("Unsupported experiment option.") + data: dict[str, Any] = { + "name": _identifier(name, "name", 400), + "id": _uuid(options.get("id", uuid4())), + "started_at": _timestamp(options.get("started_at", datetime.now(timezone.utc))), + } + if name is None: + raise ValueError("name is required.") + for key, value in options.items(): + if key in ("id", "started_at"): + continue + if key == "run_source": + if value is not None and value not in ("ci", "local", "scheduled"): + raise ValueError("run_source must be ci, local, scheduled, or None.") + data[key] = value + elif key in ("expected_item_count", "expected_result_count"): + data[key] = _count(value, key) + elif key == "dataset_revision_id": + data[key] = None if value is None else _uuid(value) + else: + data[key] = _identifier(value, key) + return _encode(data) + + +@dataclass(frozen=True) +class _Chunk: + start: int + results: tuple[EvaluationResult, ...] + body: bytes + + +def _upload_body(items: list[bytes], results: list[bytes]) -> bytes: + return ( + b'{"items":[' + b",".join(items) + b'],"results":[' + b",".join(results) + b"]}" + ) + + +def _chunks(results: tuple[EvaluationResult, ...]) -> list[_Chunk]: + # Validate the finite submission before sending any chunk. Receipts and + # prepared bodies use memory proportional to this caller-sized sequence. + identities: set[tuple[str, str]] = set() + declarations: dict[str, bytes] = {} + for result in results: + if not isinstance(result, EvaluationResult): + raise TypeError("upload_results expects EvaluationResult objects.") + key = (result.item_id, result.scorer_version_id) + if key in identities: + raise ValueError( + "Each item/scorer-version pair must occur only once in a bulk upload." + ) + identities.add(key) + if isinstance(result.item, EvaluationItem): + previous = declarations.setdefault(result.item_id, result.item._json) + if previous != result.item._json: + raise ValueError("Conflicting declarations for the same item UUID.") + + chunks: list[_Chunk] = [] + item_ids: set[str] = set() + items: list[bytes] = [] + encoded_results: list[bytes] = [] + start = 0 + size = len(_upload_body([], [])) + for index, result in enumerate(results): + # Include declarations wherever their results occur. This makes every + # chunk independently replayable even after a lost acknowledgment. + item = ( + declarations.get(result.item_id) if result.item_id not in item_ids else None + ) + extra = len(result._json) + bool(encoded_results) + if item is not None: + extra += len(item) + bool(items) + if encoded_results and ( + len(encoded_results) >= _MAX_RESULTS + or len(items) + (item is not None) > _MAX_ITEMS + or size + extra > _MAX_REQUEST_BYTES + ): + chunks.append( + _Chunk( + start, results[start:index], _upload_body(items, encoded_results) + ) + ) + start, size = index, len(_upload_body([], [])) + item_ids, items, encoded_results = set(), [], [] + item = declarations.get(result.item_id) + extra = len(result._json) + (len(item) if item is not None else 0) + if size + extra > _MAX_REQUEST_BYTES: + raise ValueError( + "One result and its item exceed the upload request byte limit." + ) + if item is not None: + item_ids.add(result.item_id) + items.append(item) + encoded_results.append(result._json) + size += extra + if encoded_results: + chunks.append( + _Chunk(start, results[start:], _upload_body(items, encoded_results)) + ) + return chunks + + +class Experiment: + """A project-bound experiment handle. Closing its client does not close the run. + + Calls return only after server acknowledgment. Complete after all intended + uploads succeed; failed uploads remain resumable with the same identities. + Operations on this handle are serialized, including all chunks of a bulk + upload. Other handles and processes remain governed by server concurrency. + """ + + def __init__( + self, + transport: SyncTransport, + project_id: int, + id: str, + *, + receipt: ExperimentReceipt | None = None, + submission: bytes | None = None, + ) -> None: + self._id = id + self.receipt = receipt + self._submission = submission + self._transport = transport + self._path = ( + f"/api/projects/{project_id}/ai_observability/offline_experiments/{id}/" + ) + self._lock = RLock() + + @property + def id(self) -> str: + """Immutable experiment UUID used by every operation on this handle.""" + return self._id + + @property + def submission(self) -> dict[str, Any] | None: + """Copy of original creation fields; unavailable on a resumed handle.""" + return None if self._submission is None else json.loads(self._submission) + + def upload_result( + self, + *, + item: EvaluationItem | str | UUID, + scorer_version_id: str | UUID, + **options: Unpack[ResultOptions], + ) -> UploadReceipt: + """Persist one result and its optional item declaration, with a receipt. + + Use the same item object for multiple scorers. A UUID reference requires + an already accepted item. ``ok`` is execution success, not score polarity. + Raises EvaluationAPIError on rejection or an unknown persistence outcome. + """ + result = EvaluationResult( + item=item, scorer_version_id=scorer_version_id, **options + ) + chunk = _chunks((result,))[0] + with self._lock: + data = self._transport.request( + "POST", self._path + "upload/", body=chunk.body + ) + return _upload_receipt(data, chunk.results) + + def upload_results(self, results: Sequence[EvaluationResult]) -> BulkUploadReceipt: + """Eagerly upload a finite sequence, chunked by count and encoded bytes. + + Returns per-chunk receipts; an empty sequence makes no requests. Memory + scales with the caller's sequence. For large datasets, call in batches. + BulkUploadError exposes committed chunks and pending_results to retry. + """ + prepared = tuple(results) + chunks = _chunks(prepared) + receipts: list[UploadReceipt] = [] + with self._lock: + for chunk in chunks: + try: + data = self._transport.request( + "POST", self._path + "upload/", body=chunk.body + ) + receipts.append(_upload_receipt(data, chunk.results)) + except EvaluationAPIError as exc: + raise BulkUploadError( + exc, + completed=BulkUploadReceipt(tuple(receipts)), + pending_results=prepared[chunk.start :], + failed_index=chunk.start, + failed_body=chunk.body, + ) from exc + return BulkUploadReceipt(tuple(receipts)) + + def _close(self, action: str) -> ExperimentReceipt: + with self._lock: + data = self._transport.request( + "POST", self._path + action + "/", body=b"{}" + ) + receipt = _experiment_receipt(data, self.id) + if receipt.status != {"complete": "completed", "fail": "failed"}[action]: + raise EvaluationAPIError(code="invalid_response", persistence="unknown") + self.receipt = receipt + return receipt + + def complete(self) -> ExperimentReceipt: + """Close as completed. Server-declared expected counts must match. + + Repeating completion is safe. A count or terminal-state conflict raises + EvaluationAPIError with the server's count/error details. + """ + return self._close("complete") + + def fail(self) -> ExperimentReceipt: + """Explicitly close as failed, retaining accepted data without count checks.""" + return self._close("fail") + + +class AsyncExperiment: + """Async experiment handle with the same persistence contract as Experiment. + + Cancellation propagates and never closes the experiment. Replay the same + items/results to recover an interrupted upload, including uncertain commits. + Operations on a handle are serialized; use it on its owning event loop. + """ + + def __init__( + self, + transport: AsyncTransport, + project_id: int, + id: str, + *, + receipt: ExperimentReceipt | None = None, + submission: bytes | None = None, + ) -> None: + self._id = id + self.receipt = receipt + self._submission = submission + self._transport = transport + self._path = ( + f"/api/projects/{project_id}/ai_observability/offline_experiments/{id}/" + ) + self._lock = asyncio.Lock() + + @property + def id(self) -> str: + """Immutable experiment UUID used by every operation on this handle.""" + return self._id + + @property + def submission(self) -> dict[str, Any] | None: + """Copy of original creation fields; unavailable on a resumed handle.""" + return None if self._submission is None else json.loads(self._submission) + + async def upload_result( + self, + *, + item: EvaluationItem | str | UUID, + scorer_version_id: str | UUID, + **options: Unpack[ResultOptions], + ) -> UploadReceipt: + """Await persistence of one result. See Experiment.upload_result.""" + result = EvaluationResult( + item=item, scorer_version_id=scorer_version_id, **options + ) + chunk = _chunks((result,))[0] + async with self._lock: + data = await self._transport.request( + "POST", self._path + "upload/", body=chunk.body + ) + return _upload_receipt(data, chunk.results) + + async def upload_results( + self, results: Sequence[EvaluationResult] + ) -> BulkUploadReceipt: + """Await sequential bounded uploads. See Experiment.upload_results.""" + prepared = tuple(results) + chunks = _chunks(prepared) + receipts: list[UploadReceipt] = [] + async with self._lock: + for chunk in chunks: + try: + data = await self._transport.request( + "POST", self._path + "upload/", body=chunk.body + ) + receipts.append(_upload_receipt(data, chunk.results)) + except EvaluationAPIError as exc: + raise BulkUploadError( + exc, + completed=BulkUploadReceipt(tuple(receipts)), + pending_results=prepared[chunk.start :], + failed_index=chunk.start, + failed_body=chunk.body, + ) from exc + return BulkUploadReceipt(tuple(receipts)) + + async def _close(self, action: str) -> ExperimentReceipt: + async with self._lock: + data = await self._transport.request( + "POST", self._path + action + "/", body=b"{}" + ) + receipt = _experiment_receipt(data, self.id) + if receipt.status != {"complete": "completed", "fail": "failed"}[action]: + raise EvaluationAPIError(code="invalid_response", persistence="unknown") + self.receipt = receipt + return receipt + + async def complete(self) -> ExperimentReceipt: + """Await completed closure, checking declared counts on the server.""" + return await self._close("complete") + + async def fail(self) -> ExperimentReceipt: + """Await explicit failed closure, retaining all accepted data.""" + return await self._close("fail") + + +class OfflineEvaluations: + """Upload offline evaluations and manage scorers in one project/environment. + + ``secret_key`` accepts personal or project secret keys with + offline_evaluation_ingestion:write. Scorer reads/writes require a personal + key with llm_analytics:read/write. The offline feature must be enabled. + ``host`` defaults to the US app API; known ingestion hosts are mapped to + their app hosts. Custom hosts are preserved. Requests use a 15 second + timeout and at most three retries by default. Scorer creation/version bumps + are never automatically retried. There is no background upload queue. + """ + + def __init__( + self, + *, + project_id: int, + secret_key: str, + host: str | None = None, + timeout: float = 15, + max_retries: int = 3, + ) -> None: + self._transport = SyncTransport( + project_id, secret_key, host, timeout, max_retries + ) + self._project_id = project_id + self.scorers = Scorers(self._transport, project_id) + self._path = f"/api/projects/{project_id}/ai_observability/offline_experiments/" + + def create_experiment( + self, *, name: str, **options: Unpack[ExperimentOptions] + ) -> Experiment: + """Create and acknowledge one execution; generate UUID/start time once. + + To resume creation, reuse id, started_at and every original field. + EvaluationAPIError.submission includes the complete prepared request. + Save explicit fields before calling if recovery from process death is + required; this SDK does not persist upload state to disk. + """ + body = _experiment_body(name, options) + submission = json.loads(body) + try: + data = self._transport.request("POST", self._path, body=body) + receipt = _experiment_receipt(data, submission["id"]) + except EvaluationAPIError as exc: + exc.submission = submission + raise + return Experiment( + self._transport, + self._project_id, + receipt.id, + receipt=receipt, + submission=body, + ) + + def resume_experiment(self, id: str | UUID) -> Experiment: + """Return a local handle to an existing UUID, without a read request. + + Existence and permissions are checked by the next server operation. + This does not recreate an experiment whose creation was never accepted. + """ + return Experiment(self._transport, self._project_id, _uuid(id)) + + def close(self) -> None: + """Release HTTP connections. Experiments retain their server state.""" + self._transport.close() + + def __enter__(self) -> OfflineEvaluations: + return self + + def __exit__( + self, + exc_type: type[BaseException] | None, + exc: BaseException | None, + traceback: TracebackType | None, + ) -> None: + self.close() + + +class AsyncOfflineEvaluations: + """Native async counterpart to OfflineEvaluations; requires posthog[async]. + + Configuration, receipts and errors match the sync client. Use ``async with`` + or ``await aclose()`` to release connections. Cancellation does not close an + experiment. Use each client on a single event loop. + """ + + def __init__( + self, + *, + project_id: int, + secret_key: str, + host: str | None = None, + timeout: float = 15, + max_retries: int = 3, + ) -> None: + self._transport = AsyncTransport( + project_id, secret_key, host, timeout, max_retries + ) + self._project_id = project_id + self.scorers = AsyncScorers(self._transport, project_id) + self._path = f"/api/projects/{project_id}/ai_observability/offline_experiments/" + + async def create_experiment( + self, *, name: str, **options: Unpack[ExperimentOptions] + ) -> AsyncExperiment: + """Await creation. See OfflineEvaluations.create_experiment for replay.""" + body = _experiment_body(name, options) + submission = json.loads(body) + try: + data = await self._transport.request("POST", self._path, body=body) + receipt = _experiment_receipt(data, submission["id"]) + except EvaluationAPIError as exc: + exc.submission = submission + raise + return AsyncExperiment( + self._transport, + self._project_id, + receipt.id, + receipt=receipt, + submission=body, + ) + + def resume_experiment(self, id: str | UUID) -> AsyncExperiment: + """Construct a handle locally; no await or read permission is needed.""" + return AsyncExperiment(self._transport, self._project_id, _uuid(id)) + + async def aclose(self) -> None: + """Release async HTTP connections without changing experiment state.""" + await self._transport.aclose() + + async def __aenter__(self) -> AsyncOfflineEvaluations: + return self + + async def __aexit__( + self, + exc_type: type[BaseException] | None, + exc: BaseException | None, + traceback: TracebackType | None, + ) -> None: + await self.aclose() diff --git a/posthog/ai/evaluations/_errors.py b/posthog/ai/evaluations/_errors.py new file mode 100644 index 000000000..2dc5ffb64 --- /dev/null +++ b/posthog/ai/evaluations/_errors.py @@ -0,0 +1,43 @@ +from __future__ import annotations + +from copy import deepcopy +from typing import Any, Literal + + +class EvaluationAPIError(Exception): + """An evaluation API request did not produce a confirmed acknowledgment. + + ``persistence`` is ``"rejected"`` for a definite rejection, ``"unknown"`` + when an attempt may have persisted, or ``"not_sent"`` before sending. + The structured response retains validation paths, conflict counts, and other + server details. Inspect these attributes explicitly: exception messages omit + response content to avoid logging evaluation data or credentials. + + Experiment creation errors expose the prepared declaration in ``submission`` + so callers can retry using the same identity and content. + """ + + def __init__( + self, + *, + status: int | None = None, + code: str | None = None, + detail: Any = None, + attr: str | None = None, + errors: Any = None, + response: dict[str, Any] | None = None, + retry_after: float | None = None, + persistence: Literal["rejected", "unknown", "not_sent"] = "unknown", + submission: dict[str, Any] | None = None, + ) -> None: + self.status = status + self.code = code + self.detail = deepcopy(detail) + self.attr = attr + self.errors = deepcopy(errors) + self.response = deepcopy(response) + self.retry_after = retry_after + self.persistence = persistence + self.submission = deepcopy(submission) + http_status = f"HTTP {status}" if status is not None else "no acknowledgment" + super().__init__(f"Offline evaluation request failed ({http_status}).") diff --git a/posthog/ai/evaluations/_scorer_types.py b/posthog/ai/evaluations/_scorer_types.py new file mode 100644 index 000000000..6dd204541 --- /dev/null +++ b/posthog/ai/evaluations/_scorer_types.py @@ -0,0 +1,145 @@ +"""Typed scorer configurations and immutable scorer metadata.""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime +from typing import Any, Generic, Literal, TypeVar + +from typing_extensions import Required, TypedDict + + +ScorerKind = Literal["boolean", "numeric", "categorical"] +ScorerOrder = Literal[ + "name", + "-name", + "kind", + "-kind", + "created_at", + "-created_at", + "updated_at", + "-updated_at", + "current_version", + "-current_version", +] + + +class BooleanScorerConfig(TypedDict, total=False): + """Boolean labels and polarity, pinned to each scorer version. + + ``true_is_failure=True`` means false passes. Omitted, false, or null means + true passes. Labels change presentation without changing stored values. + """ + + true_is_failure: bool | None + true_label: str + false_label: str + + +class NumericPassingRule(TypedDict): + """Pass at or above (``gte``), or at or below (``lte``), the threshold.""" + + operator: Literal["gte", "lte"] + threshold: float + + +class NumericScorerConfig(TypedDict, total=False): + """Inclusive bounds, input step, and optional numeric passing rule. + + An omitted or null passing rule leaves the score neutral. A configured + threshold must be finite and within the configured bounds. + """ + + min: float | None + max: float | None + step: float | None + passing_rule: NumericPassingRule | None + + +class CategoricalScorerOption(TypedDict): + """Stable category key and its human-readable label.""" + + key: str + label: str + + +class CategoricalPassingRule(TypedDict): + """Every returned category must be in ``categories`` for a score to pass.""" + + categories: list[str] + + +class CategoricalScorerConfig(TypedDict, total=False): + """Ordered options, selection constraints, and optional passing categories. + + ``options`` is required. Selection defaults to ``single``; minimum and + maximum selections only apply to ``multiple``. An omitted or null passing + rule leaves the score neutral. An empty passing list is supported only for + multiple selection and makes every accepted offline result fail. + """ + + options: Required[list[CategoricalScorerOption]] + selection_mode: Literal["single", "multiple"] + min_selections: int | None + max_selections: int | None + passing_rule: CategoricalPassingRule | None + + +ScorerConfig = BooleanScorerConfig | NumericScorerConfig | CategoricalScorerConfig +_ConfigT = TypeVar("_ConfigT", bound=ScorerConfig, covariant=True) + + +@dataclass(frozen=True) +class Scorer(Generic[_ConfigT]): + """Scorer metadata and the configuration of its current immutable version. + + Read ``current_version_id`` once to pin the version for an experiment. + Names are not unique; use ``id`` for subsequent management operations. + Configuration dictionaries are local snapshots, not live server state. + """ + + id: str + name: str + description: str + kind: ScorerKind + archived: bool + current_version: int + current_version_id: str + config: _ConfigT + created_at: datetime + updated_at: datetime + created_by: dict[str, Any] | None + team: int + + +@dataclass(frozen=True) +class ScorerVersion(Generic[_ConfigT]): + """An exact immutable scorer version, including its original polarity.""" + + id: str + definition_id: str + version: int + kind: ScorerKind + config: _ConfigT + created_at: datetime + created_by: dict[str, Any] | None + + +@dataclass(frozen=True) +class ScorerPage: + """One scorer page; pass ``next_offset`` to ``scorers.list`` to continue.""" + + count: int + results: tuple[Scorer[ScorerConfig], ...] + next: str | None + previous: str | None + next_offset: int | None + + +@dataclass(frozen=True) +class ScorerVersionPage: + """One version page; pass ``next_cursor`` to ``list_versions`` to continue.""" + + count: int + results: tuple[ScorerVersion[ScorerConfig], ...] + next_cursor: str | None diff --git a/posthog/ai/evaluations/_scorers.py b/posthog/ai/evaluations/_scorers.py new file mode 100644 index 000000000..2c5dc3d10 --- /dev/null +++ b/posthog/ai/evaluations/_scorers.py @@ -0,0 +1,572 @@ +"""Scorer management, with identical synchronous and asynchronous interfaces.""" + +from __future__ import annotations + +import copy +import json +from datetime import datetime +from typing import Any, Callable, Literal, TypeVar, cast, overload +from urllib.parse import parse_qs, urlsplit +from uuid import UUID + +from ._errors import EvaluationAPIError +from ._scorer_types import ( + BooleanScorerConfig, + CategoricalScorerConfig, + NumericScorerConfig, + Scorer, + ScorerConfig, + ScorerKind, + ScorerOrder, + ScorerPage, + ScorerVersion, + ScorerVersionPage, +) +from ._transport import AsyncTransport, SyncTransport + + +class _Omitted: + def __repr__(self) -> str: + return "" + + +_OMITTED = _Omitted() +_ResultT = TypeVar("_ResultT") +_ConfigT = TypeVar("_ConfigT", bound=ScorerConfig) + + +def _uuid(value: str | UUID) -> str: + if isinstance(value, UUID): + return str(value) + if not isinstance(value, str): + raise ValueError("Scorer and version IDs must be UUIDs.") + try: + return str(UUID(value)) + except ValueError: + raise ValueError("Scorer and version IDs must be UUIDs.") from None + + +def _integer(value: Any, *, minimum: int = 0) -> int: + if type(value) is not int or value < minimum: + raise ValueError("Expected an integer within the supported range.") + return value + + +def _text(value: Any) -> str: + if not isinstance(value, str): + raise ValueError("Expected a string.") + return value + + +def _optional_text(value: Any) -> str | None: + return None if value is None else _text(value) + + +def _kind(value: Any) -> ScorerKind: + if value not in ("boolean", "numeric", "categorical"): + raise ValueError("Unsupported scorer kind.") + return cast(ScorerKind, value) + + +def _timestamp(value: Any) -> datetime: + parsed = datetime.fromisoformat(_text(value).replace("Z", "+00:00")) + if parsed.tzinfo is None or parsed.utcoffset() is None: + raise ValueError("Expected a timestamp with a timezone.") + return parsed + + +def _mapping(value: Any) -> dict[str, Any]: + if not isinstance(value, dict): + raise ValueError("Expected an object.") + return copy.deepcopy(value) + + +def _scorer(data: dict[str, Any]) -> Scorer[ScorerConfig]: + archived = data["archived"] + if not isinstance(archived, bool): + raise ValueError("Expected boolean archived state.") + version_id = data["current_version_id"] + created_by = data["created_by"] + return Scorer( + id=_uuid(data["id"]), + name=_text(data["name"]), + description=_text(data["description"]), + kind=_kind(data["kind"]), + archived=archived, + current_version=_integer(data["current_version"], minimum=1), + current_version_id=_uuid(version_id), + config=cast(ScorerConfig, _mapping(data["config"])), + created_at=_timestamp(data["created_at"]), + updated_at=_timestamp(data["updated_at"]), + created_by=None if created_by is None else _mapping(created_by), + team=_integer(data["team"], minimum=1), + ) + + +def _version(data: dict[str, Any]) -> ScorerVersion[ScorerConfig]: + created_by = data["created_by"] + return ScorerVersion( + id=_uuid(data["id"]), + definition_id=_uuid(data["definition_id"]), + version=_integer(data["version"], minimum=1), + kind=_kind(data["kind"]), + config=cast(ScorerConfig, _mapping(data["config"])), + created_at=_timestamp(data["created_at"]), + created_by=None if created_by is None else _mapping(created_by), + ) + + +def _results(data: dict[str, Any]) -> list[dict[str, Any]]: + results = data["results"] + if not isinstance(results, list) or any( + not isinstance(row, dict) for row in results + ): + raise ValueError("Expected a list of objects.") + return results + + +def _scorer_page(data: dict[str, Any]) -> ScorerPage: + next_url = _optional_text(data["next"]) + next_offset = None + if next_url is not None: + offsets = parse_qs(urlsplit(next_url).query).get("offset") + if offsets is None or len(offsets) != 1: + raise ValueError("Missing scorer page offset.") + next_offset = _integer(int(offsets[0])) + return ScorerPage( + count=_integer(data["count"]), + results=tuple(_scorer(row) for row in _results(data)), + next=next_url, + previous=_optional_text(data["previous"]), + next_offset=next_offset, + ) + + +def _version_page(data: dict[str, Any]) -> ScorerVersionPage: + return ScorerVersionPage( + count=_integer(data["count"]), + results=tuple(_version(row) for row in _results(data)), + next_cursor=_optional_text(data["next_cursor"]), + ) + + +def _parse( + parser: Callable[[dict[str, Any]], _ResultT], data: dict[str, Any] +) -> _ResultT: + try: + return parser(data) + except (KeyError, TypeError, ValueError, OverflowError): + raise EvaluationAPIError( + code="invalid_response", + detail="The scorer API returned an invalid response.", + response=data, + persistence="unknown", + ) from None + + +def _encode(data: dict[str, Any]) -> bytes: + # Preparing bytes before transport I/O snapshots nested caller-owned values. + return json.dumps( + data, ensure_ascii=False, allow_nan=False, separators=(",", ":") + ).encode("utf-8") + + +def _create_body( + name: str, + kind: ScorerKind, + config: ScorerConfig, + description: str | None | _Omitted, +) -> bytes: + data: dict[str, Any] = { + "name": name, + "kind": _kind(kind), + "config": _mapping(config), + } + if not isinstance(description, _Omitted): + data["description"] = description + return _encode(data) + + +def _metadata( + name: str | _Omitted, + description: str | None | _Omitted, + archived: bool | _Omitted = _OMITTED, +) -> dict[str, Any]: + return { + key: value + for key, value in { + "name": name, + "description": description, + "archived": archived, + }.items() + if not isinstance(value, _Omitted) + } + + +def _version_body( + config: ScorerConfig, + base_version: int | None, + name: str | _Omitted, + description: str | None | _Omitted, +) -> bytes: + data = _metadata(name, description) + data["config"] = _mapping(config) + if base_version is not None: + data["base_version"] = _integer(base_version, minimum=1) + return _encode(data) + + +def _list_params( + limit: int, + offset: int, + search: str | None, + kind: ScorerKind | None, + archived: bool | None, + order_by: ScorerOrder | None, +) -> dict[str, str | int]: + params: dict[str, str | int] = { + "limit": _integer(limit, minimum=1), + "offset": _integer(offset), + } + if search is not None: + params["search"] = search + if kind is not None: + params["kind"] = _kind(kind) + if archived is not None: + if not isinstance(archived, bool): + raise ValueError("archived must be a boolean.") + params["archived"] = "true" if archived else "false" + if order_by is not None: + params["order_by"] = order_by + return params + + +def _version_params(limit: int, cursor: str | None) -> dict[str, str | int]: + if _integer(limit, minimum=1) > 100: + raise ValueError("Version page limit must be between 1 and 100.") + params: dict[str, str | int] = {"limit": limit} + if cursor is not None: + params["cursor"] = cursor + return params + + +class Scorers: + """Manage scorers using a personal key with ``llm_analytics`` scopes. + + Access through ``OfflineEvaluations.scorers``. Creating a scorer or version + is not idempotent: mutations are never retried automatically. If a response + is lost, inspect server state before repeating the operation. + """ + + def __init__(self, transport: SyncTransport, project_id: int) -> None: + self._transport = transport + self._base = f"/api/projects/{project_id}/llm_analytics/score_definitions/" + + def _path(self, scorer_id: str | UUID) -> str: + return f"{self._base}{_uuid(scorer_id)}/" + + @overload + def create( + self, + *, + name: str, + kind: Literal["boolean"], + config: BooleanScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[BooleanScorerConfig]: ... + + @overload + def create( + self, + *, + name: str, + kind: Literal["numeric"], + config: NumericScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[NumericScorerConfig]: ... + + @overload + def create( + self, + *, + name: str, + kind: Literal["categorical"], + config: CategoricalScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[CategoricalScorerConfig]: ... + + def create( + self, + *, + name: str, + kind: ScorerKind, + config: ScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[ScorerConfig]: + """Create an active scorer and its first immutable configuration version.""" + return _parse( + _scorer, + self._transport.request( + "POST", + self._base, + body=_create_body(name, kind, config, description), + retry_safe=False, + ), + ) + + def get(self, scorer_id: str | UUID) -> Scorer[ScorerConfig]: + """Fetch a scorer by its stable UUID, including ``current_version_id``.""" + return _parse(_scorer, self._transport.request("GET", self._path(scorer_id))) + + def list( + self, + *, + limit: int = 100, + offset: int = 0, + search: str | None = None, + kind: ScorerKind | None = None, + archived: bool | None = None, + order_by: ScorerOrder | None = None, + ) -> ScorerPage: + """Fetch one page. Omit ``archived`` to list active scorers. + + Keep the same filters and pass the returned ``next_offset`` to continue. + Search matches names and descriptions; scorer names are not unique. + """ + return _parse( + _scorer_page, + self._transport.request( + "GET", + self._base, + params=_list_params(limit, offset, search, kind, archived, order_by), + ), + ) + + def update( + self, + scorer_id: str | UUID, + *, + name: str | _Omitted = _OMITTED, + description: str | None | _Omitted = _OMITTED, + archived: bool | _Omitted = _OMITTED, + ) -> Scorer[ScorerConfig]: + """Update metadata. Use ``create_version`` to change configuration. + + Omitted fields are preserved. An explicit null description clears it. + """ + return _parse( + _scorer, + self._transport.request( + "PATCH", + self._path(scorer_id), + body=_encode(_metadata(name, description, archived)), + retry_safe=False, + ), + ) + + def create_version( + self, + scorer_id: str | UUID, + *, + config: _ConfigT, + base_version: int | None = None, + name: str | _Omitted = _OMITTED, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[_ConfigT]: + """Create an immutable version and return the updated scorer. + + Pass the observed ``current_version`` as ``base_version`` to detect + concurrent writes. Conflicts are surfaced without automatic refresh. + Identical configuration is allowed, for example after a prompt change. + """ + return cast( + Scorer[_ConfigT], + _parse( + _scorer, + self._transport.request( + "POST", + self._path(scorer_id) + "new_version/", + body=_version_body(config, base_version, name, description), + retry_safe=False, + ), + ), + ) + + def get_version( + self, scorer_id: str | UUID, version_id: str | UUID + ) -> ScorerVersion[ScorerConfig]: + """Fetch an exact version UUID belonging to the specified scorer.""" + return _parse( + _version, + self._transport.request( + "GET", self._path(scorer_id) + f"versions/{_uuid(version_id)}/" + ), + ) + + def list_versions( + self, scorer_id: str | UUID, *, limit: int = 50, cursor: str | None = None + ) -> ScorerVersionPage: + """Fetch versions newest first, using the returned cursor to continue.""" + return _parse( + _version_page, + self._transport.request( + "GET", + self._path(scorer_id) + "versions/", + params=_version_params(limit, cursor), + ), + ) + + +class AsyncScorers: + """Async counterpart of :class:`Scorers`, exposed by AsyncOfflineEvaluations.""" + + def __init__(self, transport: AsyncTransport, project_id: int) -> None: + self._transport = transport + self._base = f"/api/projects/{project_id}/llm_analytics/score_definitions/" + + def _path(self, scorer_id: str | UUID) -> str: + return f"{self._base}{_uuid(scorer_id)}/" + + @overload + async def create( + self, + *, + name: str, + kind: Literal["boolean"], + config: BooleanScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[BooleanScorerConfig]: ... + + @overload + async def create( + self, + *, + name: str, + kind: Literal["numeric"], + config: NumericScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[NumericScorerConfig]: ... + + @overload + async def create( + self, + *, + name: str, + kind: Literal["categorical"], + config: CategoricalScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[CategoricalScorerConfig]: ... + + async def create( + self, + *, + name: str, + kind: ScorerKind, + config: ScorerConfig, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[ScorerConfig]: + """Create an active scorer and its first version; never retry mutations.""" + return _parse( + _scorer, + await self._transport.request( + "POST", + self._base, + body=_create_body(name, kind, config, description), + retry_safe=False, + ), + ) + + async def get(self, scorer_id: str | UUID) -> Scorer[ScorerConfig]: + """Fetch a scorer by stable UUID, including its current version UUID.""" + return _parse( + _scorer, await self._transport.request("GET", self._path(scorer_id)) + ) + + async def list( + self, + *, + limit: int = 100, + offset: int = 0, + search: str | None = None, + kind: ScorerKind | None = None, + archived: bool | None = None, + order_by: ScorerOrder | None = None, + ) -> ScorerPage: + """Fetch one page; preserve filters and use ``next_offset`` to continue.""" + return _parse( + _scorer_page, + await self._transport.request( + "GET", + self._base, + params=_list_params(limit, offset, search, kind, archived, order_by), + ), + ) + + async def update( + self, + scorer_id: str | UUID, + *, + name: str | _Omitted = _OMITTED, + description: str | None | _Omitted = _OMITTED, + archived: bool | _Omitted = _OMITTED, + ) -> Scorer[ScorerConfig]: + """Update metadata; omission preserves values and null clears description.""" + return _parse( + _scorer, + await self._transport.request( + "PATCH", + self._path(scorer_id), + body=_encode(_metadata(name, description, archived)), + retry_safe=False, + ), + ) + + async def create_version( + self, + scorer_id: str | UUID, + *, + config: _ConfigT, + base_version: int | None = None, + name: str | _Omitted = _OMITTED, + description: str | None | _Omitted = _OMITTED, + ) -> Scorer[_ConfigT]: + """Create a version; ``base_version`` detects concurrent configuration edits. + + Returns the updated scorer. Identical configuration is allowed. + """ + return cast( + Scorer[_ConfigT], + _parse( + _scorer, + await self._transport.request( + "POST", + self._path(scorer_id) + "new_version/", + body=_version_body(config, base_version, name, description), + retry_safe=False, + ), + ), + ) + + async def get_version( + self, scorer_id: str | UUID, version_id: str | UUID + ) -> ScorerVersion[ScorerConfig]: + """Fetch an exact immutable version by scorer and version UUIDs.""" + return _parse( + _version, + await self._transport.request( + "GET", self._path(scorer_id) + f"versions/{_uuid(version_id)}/" + ), + ) + + async def list_versions( + self, scorer_id: str | UUID, *, limit: int = 50, cursor: str | None = None + ) -> ScorerVersionPage: + """Fetch versions newest first, using the returned cursor to continue.""" + return _parse( + _version_page, + await self._transport.request( + "GET", + self._path(scorer_id) + "versions/", + params=_version_params(limit, cursor), + ), + ) diff --git a/posthog/ai/evaluations/_serialization.py b/posthog/ai/evaluations/_serialization.py new file mode 100644 index 000000000..b0d8bf448 --- /dev/null +++ b/posthog/ai/evaluations/_serialization.py @@ -0,0 +1,184 @@ +from __future__ import annotations + +import json +import math +from datetime import datetime, timezone +from typing import Any, Iterable, Union, cast +from uuid import UUID + +JSONValue = Union[ + None, bool, int, float, str, list["JSONValue"], dict[str, "JSONValue"] +] + + +class _Unset: + def __repr__(self) -> str: + return "" + + +_UNSET = _Unset() +_MAX_ITEMS = 1000 +_MAX_RESULTS = 1000 +_MAX_REQUEST_BYTES = 5 * 1024 * 1024 +_MAX_ITEM_BYTES = 1024 * 1024 +_MAX_RESULT_BYTES = 256 * 1024 + + +def _uuid(value: object) -> str: + if not isinstance(value, (str, UUID)): + raise ValueError("Provide a UUID or a valid UUID string.") + try: + return str(UUID(str(value))) + except ValueError as exc: + raise ValueError("Provide a valid UUID.") from exc + + +def _timestamp(value: datetime | str) -> str: + if isinstance(value, str): + value = datetime.fromisoformat(value.replace("Z", "+00:00")) + if not isinstance(value, datetime) or value.utcoffset() is None: + raise ValueError("Timestamps must include a timezone.") + return value.astimezone(timezone.utc).isoformat() + + +def _identifier(value: object, field: str, limit: int = 255) -> str | None: + if value is not None and ( + not isinstance(value, str) or not value or len(value) > limit or "\x00" in value + ): + raise ValueError( + f"{field} must be a nonempty string of at most {limit} characters." + ) + return value + + +def _encode(value: object) -> bytes: + try: + return json.dumps( + value, + ensure_ascii=False, + allow_nan=False, + separators=(",", ":"), + sort_keys=True, + ).encode("utf-8") + except (TypeError, ValueError, UnicodeError) as exc: + raise ValueError( + "Provide JSON values with finite numbers and valid Unicode." + ) from exc + + +def _validate_json(value: object, depth: int = 1) -> None: + if isinstance(value, (dict, list)): + if depth > 32: + raise ValueError("Payload JSON nesting must not exceed 32 levels.") + children: Iterable[object] + if isinstance(value, dict): + if any(not isinstance(key, str) or "\x00" in key for key in value): + raise ValueError( + "JSON object keys must be strings without null characters." + ) + children = value.values() + else: + children = value + for child in children: + _validate_json(child, depth + 1) + elif isinstance(value, str): + if "\x00" in value: + raise ValueError("JSON strings must not contain null characters.") + elif value is not None and type(value) not in (bool, int, float): + raise ValueError("Provide JSON values, without automatic object conversion.") + + +def _payload( + explicit: dict[str, JSONValue] | _Unset, + fields: dict[str, JSONValue | _Unset], + limit: int, +) -> dict[str, JSONValue] | _Unset: + supplied = { + key: value for key, value in fields.items() if not isinstance(value, _Unset) + } + if not isinstance(explicit, _Unset): + if supplied: + raise ValueError("Supply payload or individual payload fields, not both.") + if not isinstance(explicit, dict) or set(explicit) - set(fields): + raise ValueError( + "payload must be an object containing supported fields only." + ) + supplied = explicit + elif not supplied: + return _UNSET + metadata = supplied.get("metadata") + if metadata is not None and not isinstance(metadata, dict): + raise ValueError("metadata must be a JSON object or None.") + for field in ("reasoning", "error_message"): + if supplied.get(field) is not None and not isinstance(supplied[field], str): + raise ValueError(f"{field} must be a string or None.") + _validate_json(supplied) + encoded = _encode(supplied) + if len(encoded) > limit: + raise ValueError(f"Payload exceeds {limit} bytes of UTF-8 JSON.") + return json.loads(encoded) + + +def _score(value: object) -> bool | float | list[str]: + if type(value) is bool: + return value + if type(value) in (int, float): + try: + number = float(cast(Union[int, float], value)) + except OverflowError as exc: + raise ValueError( + "Numeric scores must fit a finite binary64 number." + ) from exc + if not math.isfinite(number): + raise ValueError("Numeric scores must be finite.") + if type(value) is int and number != value: + raise ValueError( + "Integer scores must be exactly representable as a binary64 number." + ) + return number + if ( + isinstance(value, list) + and value + and all( + isinstance(key, str) and key and len(key) <= 128 and "\x00" not in key + for key in value + ) + ): + if len(set(value)) != len(value): + raise ValueError("Categorical scores must contain distinct keys.") + return sorted(value) + raise ValueError( + "A score must be a boolean, finite number, or nonempty list of category keys." + ) + + +def _count(value: object, field: str) -> int | None: + if value is not None and (type(value) is not int or not 0 <= value <= 2**63 - 1): + raise ValueError( + f"{field} must be a nonnegative signed 64-bit integer or None." + ) + return value + + +def _object(data: Any) -> dict[str, Any]: + if not isinstance(data, dict): + raise ValueError("Expected a response object.") + return data + + +def _boolean(value: Any) -> bool: + if type(value) is not bool: + raise ValueError("Expected a boolean.") + return value + + +def _response_count(value: Any) -> int: + if type(value) is not int or value < 0: + raise ValueError("Expected a nonnegative count.") + return value + + +def _date(value: Any) -> datetime: + if not isinstance(value, str): + raise ValueError("Expected an ISO timestamp.") + return datetime.fromisoformat(_timestamp(value)) diff --git a/posthog/ai/evaluations/_transport.py b/posthog/ai/evaluations/_transport.py new file mode 100644 index 000000000..f6489c502 --- /dev/null +++ b/posthog/ai/evaluations/_transport.py @@ -0,0 +1,308 @@ +from __future__ import annotations + +import asyncio +import importlib +import math +import random +import time +from datetime import datetime, timezone +from email.utils import parsedate_to_datetime +from typing import Any +from urllib.parse import urlsplit, urlunsplit + +import requests + +from ...version import VERSION +from ._errors import EvaluationAPIError + +_MAX_RETRY_DELAY = 30.0 +_RETRY_STATUSES = {408, 429, 500, 502, 503, 504} +_CLOUD_HOSTS = { + "app.posthog.com": "us.posthog.com", + "us.i.posthog.com": "us.posthog.com", + "eu.i.posthog.com": "eu.posthog.com", +} + + +def _normalize_host(host: str | None) -> str: + parsed = urlsplit((host or "https://us.posthog.com").strip()) + if ( + parsed.scheme not in {"http", "https"} + or not parsed.hostname + or parsed.username is not None + or parsed.password is not None + or parsed.query + or parsed.fragment + ): + raise ValueError("host must be an absolute HTTP(S) URL without credentials") + # Accessing port also validates malformed port numbers before any I/O. + netloc = parsed.netloc + if parsed.port in {None, 443} and parsed.scheme == "https": + netloc = _CLOUD_HOSTS.get(parsed.hostname, netloc) + return urlunsplit((parsed.scheme, netloc, parsed.path.rstrip("/"), "", "")) + + +def _validate_options( + project_id: int, secret_key: str, timeout: float, max_retries: int +) -> None: + if ( + isinstance(project_id, bool) + or not isinstance(project_id, int) + or project_id <= 0 + ): + raise ValueError("project_id must be a positive integer") + if ( + not isinstance(secret_key, str) + or not secret_key + or any(character.isspace() for character in secret_key) + ): + raise ValueError("secret_key must be a nonempty credential without whitespace") + if ( + isinstance(timeout, bool) + or not isinstance(timeout, (int, float)) + or not math.isfinite(timeout) + or timeout <= 0 + ): + raise ValueError("timeout must be a finite positive number of seconds") + if ( + isinstance(max_retries, bool) + or not isinstance(max_retries, int) + or max_retries < 0 + ): + raise ValueError("max_retries must be a nonnegative integer") + + +def _request_url(host: str, path: str) -> str: + if not path.startswith("/") or path.startswith("//") or "?" in path or "#" in path: + raise ValueError( + "path must be an absolute API path without a query or fragment" + ) + return host + path + + +def _retry_after(value: str | None) -> float | None: + if value is None: + return None + try: + seconds = float(value) + except ValueError: + try: + retry_at = parsedate_to_datetime(value) + if retry_at.tzinfo is None: + retry_at = retry_at.replace(tzinfo=timezone.utc) + seconds = (retry_at - datetime.now(timezone.utc)).total_seconds() + except (ValueError, TypeError, OverflowError): + return None + return max(0.0, seconds) if math.isfinite(seconds) else None + + +def _retry_delay( + error: EvaluationAPIError, attempt: int, max_retries: int, retry_safe: bool +) -> float | None: + if ( + not retry_safe + or attempt >= max_retries + or (error.status is not None and error.status not in _RETRY_STATUSES) + or error.code == "invalid_response" + ): + return None + # A longer Retry-After must surface to the caller. Clamping it and retrying + # early would break the server's rate limit contract. + if error.retry_after is not None and error.retry_after > _MAX_RETRY_DELAY: + return None + backoff = min(_MAX_RETRY_DELAY, 0.5 * 2 ** min(attempt, 6)) + backoff = min(_MAX_RETRY_DELAY, backoff * random.uniform(0.8, 1.2)) + return max(backoff, error.retry_after or 0.0) + + +def _parse_response( + status: int, payload: Any, retry_after: str | None +) -> dict[str, Any]: + if 200 <= status < 300: + if isinstance(payload, dict): + return payload + raise EvaluationAPIError( + status=status, + code="invalid_response", + detail="Expected a JSON object in the server acknowledgment.", + ) + response = payload if isinstance(payload, dict) else None + fields = response or {} + raise EvaluationAPIError( + status=status, + code=fields.get("code"), + detail=fields.get("detail", payload), + attr=fields.get("attr"), + errors=fields.get("errors"), + response=response, + retry_after=_retry_after(retry_after), + persistence="rejected" if 400 <= status < 500 and status != 408 else "unknown", + ) + + +class _BearerAuth(requests.auth.AuthBase): + def __init__(self, secret_key: str) -> None: + self._secret_key = secret_key + + def __call__(self, request: requests.PreparedRequest) -> requests.PreparedRequest: + # Setting auth explicitly prevents requests from substituting .netrc + # credentials for this client's API credential. + request.headers["Authorization"] = f"Bearer {self._secret_key}" + return request + + +class SyncTransport: + def __init__( + self, + project_id: int, + secret_key: str, + host: str | None = None, + timeout: float = 15, + max_retries: int = 3, + ) -> None: + _validate_options(project_id, secret_key, timeout, max_retries) + self._host = _normalize_host(host) + self._timeout = timeout + self._max_retries = max_retries + self._closed = False + self._session = requests.Session() + self._session.auth = _BearerAuth(secret_key) + self._session.headers.update( + { + "Accept": "application/json", + "Content-Type": "application/json", + "User-Agent": f"posthog-python/{VERSION}", + } + ) + self._session.mount("http://", requests.adapters.HTTPAdapter(max_retries=0)) + self._session.mount("https://", requests.adapters.HTTPAdapter(max_retries=0)) + + def request( + self, + method: str, + path: str, + *, + body: bytes | None = None, + params: dict[str, str | int] | None = None, + retry_safe: bool = True, + ) -> dict[str, Any]: + if self._closed: + raise EvaluationAPIError(code="client_closed", persistence="not_sent") + url = _request_url(self._host, path) + had_unknown_attempt = False + for attempt in range(self._max_retries + 1): + try: + response = self._session.request( + method, + url, + data=body, + params=params, + timeout=self._timeout, + allow_redirects=False, + ) + try: + try: + payload = response.json() + except ValueError: + payload = None + return _parse_response( + response.status_code, + payload, + response.headers.get("Retry-After"), + ) + finally: + response.close() + except requests.exceptions.RequestException: + error = EvaluationAPIError(code="transport_error") + except EvaluationAPIError as exc: + error = exc + had_unknown_attempt = had_unknown_attempt or error.persistence == "unknown" + if had_unknown_attempt: + error.persistence = "unknown" + delay = _retry_delay(error, attempt, self._max_retries, retry_safe) + if delay is None: + raise error from None + time.sleep(delay) + raise AssertionError("unreachable") + + def close(self) -> None: + self._closed = True + self._session.close() + + +class AsyncTransport: + def __init__( + self, + project_id: int, + secret_key: str, + host: str | None = None, + timeout: float = 15, + max_retries: int = 3, + ) -> None: + _validate_options(project_id, secret_key, timeout, max_retries) + self._host = _normalize_host(host) + try: + self._httpx = importlib.import_module("httpx") + except ImportError: + raise RuntimeError( + "Async evaluations require httpx. Install it with `posthog[async]`." + ) from None + self._max_retries = max_retries + self._closed = False + self._client = self._httpx.AsyncClient( + headers={ + "Authorization": f"Bearer {secret_key}", + "Accept": "application/json", + "Content-Type": "application/json", + "User-Agent": f"posthog-python/{VERSION}", + }, + timeout=timeout, + follow_redirects=False, + ) + + async def request( + self, + method: str, + path: str, + *, + body: bytes | None = None, + params: dict[str, str | int] | None = None, + retry_safe: bool = True, + ) -> dict[str, Any]: + if self._closed: + raise EvaluationAPIError(code="client_closed", persistence="not_sent") + url = _request_url(self._host, path) + had_unknown_attempt = False + for attempt in range(self._max_retries + 1): + try: + response = await self._client.request( + method, url, content=body, params=params + ) + try: + try: + payload = response.json() + except ValueError: + payload = None + return _parse_response( + response.status_code, + payload, + response.headers.get("Retry-After"), + ) + finally: + await response.aclose() + except self._httpx.RequestError: + error = EvaluationAPIError(code="transport_error") + except EvaluationAPIError as exc: + error = exc + had_unknown_attempt = had_unknown_attempt or error.persistence == "unknown" + if had_unknown_attempt: + error.persistence = "unknown" + delay = _retry_delay(error, attempt, self._max_retries, retry_safe) + if delay is None: + raise error from None + await asyncio.sleep(delay) + raise AssertionError("unreachable") + + async def aclose(self) -> None: + self._closed = True + await self._client.aclose() diff --git a/posthog/ai/evaluations/_types.py b/posthog/ai/evaluations/_types.py new file mode 100644 index 000000000..c6862b37d --- /dev/null +++ b/posthog/ai/evaluations/_types.py @@ -0,0 +1,394 @@ +from __future__ import annotations + +import json +from dataclasses import dataclass, field +from datetime import datetime +from typing import Any, Literal, Mapping +from uuid import UUID, uuid4 + +from typing_extensions import TypedDict + +from ._errors import EvaluationAPIError +from ._serialization import ( + JSONValue, + _MAX_ITEM_BYTES, + _MAX_RESULT_BYTES, + _UNSET, + _Unset, + _boolean, + _date, + _encode, + _identifier, + _object, + _payload, + _response_count, + _score, + _timestamp, + _uuid, +) + +RunSource = Literal["ci", "local", "scheduled"] +ResultStatus = Literal["ok", "error", "skipped", "not_applicable"] +ScoreValue = bool | int | float | list[str] + + +class ExperimentOptions(TypedDict, total=False): + """Optional experiment creation fields. Persist these for exact replay. + + UUID and start time are generated once if omitted. ``run_source`` defaults + to unspecified. Expected counts include every outcome status. Dataset and + application identifiers describe the original execution, not upload time. + """ + + id: str | UUID + started_at: datetime | str + run_source: RunSource | None + expected_item_count: int | None + expected_result_count: int | None + suite_key: str | None + dataset_source: str | None + dataset_identifier: str | None + dataset_revision_identifier: str | None + dataset_revision_id: str | UUID | None + application_version: str | None + model_version: str | None + prompt_version: str | None + + +class ResultOptions(TypedDict, total=False): + """Score outcome, trace references, and optional result payload fields. + + ``value`` is required for ``ok`` (the default). ``error_code`` and + ``error_message`` are permitted only for ``error``. Use either ``payload`` + or individual payload fields; omitted values differ from explicit None. + """ + + value: ScoreValue | None + status: ResultStatus + error_code: str | None + reasoning: str | None + error_message: str | None + metadata: dict[str, JSONValue] | None + evaluator_trace_id: str | None + evaluated_at: datetime | str | None + payload: dict[str, JSONValue] + + +@dataclass(frozen=True, init=False) +class EvaluationItem: + """An immutable input/output execution, shared by all its scorer results. + + Construction is local and generates a UUID unless ``id`` is supplied. + Payload values are copied as JSON. Omitted fields differ from explicit None. + Use ``payload={}`` for an explicitly empty payload. Save ``to_dict()`` to + restore exactly the same declaration after restarting the process. + """ + + id: str + _json: bytes = field(repr=False) + + def __init__( + self, + *, + id: str | UUID | None = None, + input: JSONValue | _Unset = _UNSET, + output: JSONValue | _Unset = _UNSET, + expected_output: JSONValue | _Unset = _UNSET, + metadata: dict[str, JSONValue] | None | _Unset = _UNSET, + case_key: str | None = None, + trial: str | None = None, + dataset_item_identifier: str | None = None, + dataset_item_version_identifier: str | None = None, + dataset_item_version_id: str | UUID | None = None, + application_trace_id: str | None = None, + payload: dict[str, JSONValue] | _Unset = _UNSET, + ) -> None: + identity = _uuid(uuid4() if id is None else id) + data: dict[str, Any] = {"id": identity} + for key, value in { + "case_key": case_key, + "trial": trial, + "dataset_item_identifier": dataset_item_identifier, + "dataset_item_version_identifier": dataset_item_version_identifier, + "application_trace_id": application_trace_id, + }.items(): + if value is not None: + data[key] = _identifier(value, key) + if dataset_item_version_id is not None: + data["dataset_item_version_id"] = _uuid(dataset_item_version_id) + prepared = _payload( + payload, + { + "input": input, + "output": output, + "expected_output": expected_output, + "metadata": metadata, + }, + _MAX_ITEM_BYTES, + ) + if not isinstance(prepared, _Unset): + data["payload"] = prepared + object.__setattr__(self, "id", identity) + object.__setattr__(self, "_json", _encode(data)) + + def to_dict(self) -> dict[str, Any]: + """Return a defensive, JSON-serializable copy of the complete declaration.""" + return json.loads(self._json) + + @classmethod + def from_dict(cls, data: Mapping[str, Any]) -> EvaluationItem: + """Restore a saved declaration, including its required original UUID.""" + if "id" not in data or data["id"] is None: + raise ValueError("Restoring an item requires its original id.") + return cls(**dict(data)) + + +@dataclass(frozen=True, init=False) +class EvaluationResult: + """One scorer outcome for an item and an exact scorer-version UUID. + + ``ok`` means the evaluator succeeded; the scorer's pinned configuration + determines whether the value passes. Other statuses contain no score. + Pass an item object to declare it, or an item UUID already accepted by the + experiment. Payloads and values are snapshotted during construction. + """ + + item: EvaluationItem | str + scorer_version_id: str + _json: bytes = field(repr=False) + + def __init__( + self, + *, + item: EvaluationItem | str | UUID, + scorer_version_id: str | UUID, + value: ScoreValue | None = None, + status: ResultStatus = "ok", + error_code: str | None = None, + reasoning: str | None | _Unset = _UNSET, + error_message: str | None | _Unset = _UNSET, + metadata: dict[str, JSONValue] | None | _Unset = _UNSET, + evaluator_trace_id: str | None = None, + evaluated_at: datetime | str | None = None, + payload: dict[str, JSONValue] | _Unset = _UNSET, + ) -> None: + item_id = item.id if isinstance(item, EvaluationItem) else _uuid(item) + version_id = _uuid(scorer_version_id) + if status not in ("ok", "error", "skipped", "not_applicable"): + raise ValueError("Unsupported result status.") + if status != "ok" and value is not None: + raise ValueError("Only an ok result may contain a score.") + if status != "error" and error_code is not None: + raise ValueError("Only an error result may contain an error_code.") + data: dict[str, Any] = { + "item_id": item_id, + "scorer_version_id": version_id, + "status": status, + } + if status == "ok": + data["value"] = _score(value) + if error_code is not None: + data["error_code"] = _identifier(error_code, "error_code", 128) + if evaluator_trace_id is not None: + data["evaluator_trace_id"] = _identifier( + evaluator_trace_id, "evaluator_trace_id" + ) + if evaluated_at is not None: + data["evaluated_at"] = _timestamp(evaluated_at) + prepared = _payload( + payload, + { + "reasoning": reasoning, + "error_message": error_message, + "metadata": metadata, + }, + _MAX_RESULT_BYTES, + ) + if not isinstance(prepared, _Unset): + if status != "error" and prepared.get("error_message") is not None: + raise ValueError("Only an error result may contain an error_message.") + data["payload"] = prepared + object.__setattr__( + self, "item", item if isinstance(item, EvaluationItem) else item_id + ) + object.__setattr__(self, "scorer_version_id", version_id) + object.__setattr__(self, "_json", _encode(data)) + + @property + def item_id(self) -> str: + """The stable item UUID used in this result's identity.""" + return self.item.id if isinstance(self.item, EvaluationItem) else self.item + + def to_dict(self) -> dict[str, Any]: + """Return the result declaration; save the item's declaration separately.""" + return json.loads(self._json) + + @classmethod + def from_dict( + cls, data: Mapping[str, Any], *, item: EvaluationItem | None = None + ) -> EvaluationResult: + """Restore a result, optionally attaching its original reusable item.""" + fields = dict(data) + item_id = _uuid(fields.pop("item_id")) + if item is not None and item.id != item_id: + raise ValueError("The restored result and item IDs do not match.") + return cls(item=item if item is not None else item_id, **fields) + + +@dataclass(frozen=True) +class ExperimentReceipt: + """Server acknowledgment of creation or closure, including accepted counts.""" + + id: str + status: Literal["uploading", "completed", "failed"] + created: bool + started_at: datetime + created_at: datetime + finished_at: datetime | None + expected_item_count: int | None + expected_result_count: int | None + accepted_item_count: int + accepted_result_count: int + + +@dataclass(frozen=True) +class ItemReceipt: + """Persisted item identity; ``created=False`` acknowledges an exact replay.""" + + id: str + created: bool + accepted_at: datetime + + +@dataclass(frozen=True) +class ResultReceipt(ItemReceipt): + """Persisted result UUID and the item/scorer-version pair it acknowledges.""" + + item_id: str + scorer_version_id: str + + +@dataclass(frozen=True) +class UploadReceipt: + """Acknowledgment of one atomic HTTP upload, including duplicate records.""" + + items: tuple[ItemReceipt, ...] + results: tuple[ResultReceipt, ...] + + +@dataclass(frozen=True) +class BulkUploadReceipt: + """Acknowledged chunks in input order. Separate chunks commit independently.""" + + chunks: tuple[UploadReceipt, ...] + + +class BulkUploadError(EvaluationAPIError): + """A stopped bulk upload with acknowledged chunks and reusable pending work. + + Retry ``pending_results`` using the same experiment. A persistence value of + ``unknown`` means the failed chunk may have committed; exact replay is safe. + ``failed_index`` is the index of its first result in the original sequence. + """ + + def __init__( + self, + cause: EvaluationAPIError, + *, + completed: BulkUploadReceipt, + pending_results: tuple[EvaluationResult, ...], + failed_index: int, + failed_body: bytes, + ) -> None: + super().__init__( + status=cause.status, + code=cause.code, + detail=cause.detail, + attr=cause.attr, + errors=cause.errors, + response=cause.response, + retry_after=cause.retry_after, + persistence=cause.persistence, + ) + self.completed = completed + self.pending_results = pending_results + self.failed_index = failed_index + self._failed_body = failed_body + + @property + def failed_chunk(self) -> dict[str, Any]: + """A defensive copy of the exact request whose acknowledgment is missing.""" + return json.loads(self._failed_body) + + +def _experiment_receipt(data: dict[str, Any], identity: str) -> ExperimentReceipt: + try: + if _uuid(data["id"]) != identity or data["status"] not in ( + "uploading", + "completed", + "failed", + ): + raise ValueError("Mismatched experiment receipt.") + return ExperimentReceipt( + id=identity, + status=data["status"], + created=_boolean(data["created"]), + started_at=_date(data["started_at"]), + created_at=_date(data["created_at"]), + finished_at=None + if data["finished_at"] is None + else _date(data["finished_at"]), + expected_item_count=None + if data["expected_item_count"] is None + else _response_count(data["expected_item_count"]), + expected_result_count=None + if data["expected_result_count"] is None + else _response_count(data["expected_result_count"]), + accepted_item_count=_response_count(data["accepted_item_count"]), + accepted_result_count=_response_count(data["accepted_result_count"]), + ) + except (KeyError, TypeError, ValueError): + raise EvaluationAPIError( + code="invalid_response", persistence="unknown" + ) from None + + +def _upload_receipt( + data: dict[str, Any], results: tuple[EvaluationResult, ...] +) -> UploadReceipt: + try: + if not isinstance(data["items"], list) or not isinstance(data["results"], list): + raise ValueError("Expected acknowledgment lists.") + items = tuple( + ItemReceipt( + id=_uuid(_object(row)["id"]), + created=_boolean(row["created"]), + accepted_at=_date(row["accepted_at"]), + ) + for row in data["items"] + ) + outcomes = tuple( + ResultReceipt( + id=_uuid(_object(row)["id"]), + created=_boolean(row["created"]), + accepted_at=_date(row["accepted_at"]), + item_id=_uuid(row["item_id"]), + scorer_version_id=_uuid(row["scorer_version_id"]), + ) + for row in data["results"] + ) + if len(items) != len({item.id for item in items}) or { + item.id for item in items + } != {result.item_id for result in results}: + raise ValueError("Incomplete item acknowledgments.") + if [(result.item_id, result.scorer_version_id) for result in outcomes] != [ + (result.item_id, result.scorer_version_id) for result in results + ]: + raise ValueError("Incomplete result acknowledgments.") + if len({result.id for result in outcomes}) != len(outcomes): + raise ValueError("Duplicate result identities in acknowledgment.") + return UploadReceipt(items=items, results=outcomes) + except (KeyError, TypeError, ValueError): + raise EvaluationAPIError( + code="invalid_response", persistence="unknown" + ) from None diff --git a/posthog/test/ai/test_evaluations.py b/posthog/test/ai/test_evaluations.py new file mode 100644 index 000000000..f6f2f96f9 --- /dev/null +++ b/posthog/test/ai/test_evaluations.py @@ -0,0 +1,611 @@ +"""Offline evaluation contracts exercised through the public sync and async API.""" + +import asyncio +import importlib.util +import inspect +import json +from datetime import datetime, timezone +from pathlib import Path +from unittest.mock import AsyncMock, Mock +from uuid import UUID, uuid4, uuid5 + +import pytest + +from posthog.ai.evaluations import ( + AsyncOfflineEvaluations, + BulkUploadError, + EvaluationAPIError, + EvaluationItem, + EvaluationResult, + OfflineEvaluations, +) +from posthog.ai.evaluations import _client +from posthog.ai.evaluations._transport import AsyncTransport, SyncTransport + + +EXPERIMENT_ID = "8cba6d9e-d2ee-4b23-a545-a61057da5870" +VERSION_ID = "f72481ec-1310-41de-b5bb-d20e1c59f121" +OTHER_VERSION_ID = "f72481ec-1310-41de-b5bb-d20e1c59f122" +ACCEPTED_AT = "2026-10-02T12:00:00Z" +BASE = "/api/projects/123/ai_observability/offline_experiments/" + + +def experiment_response(id=EXPERIMENT_ID, **overrides): + return { + "id": id, + "status": "uploading", + "created": True, + "started_at": ACCEPTED_AT, + "created_at": ACCEPTED_AT, + "finished_at": None, + "expected_item_count": None, + "expected_result_count": None, + "accepted_item_count": 0, + "accepted_result_count": 0, + **overrides, + } + + +def upload_response(body, *, created=True): + data = json.loads(body) + return { + "items": [ + {"id": id, "created": created, "accepted_at": ACCEPTED_AT} + for id in dict.fromkeys(result["item_id"] for result in data["results"]) + ], + "results": [ + { + "id": str(uuid5(UUID(result["item_id"]), result["scorer_version_id"])), + "item_id": result["item_id"], + "scorer_version_id": result["scorer_version_id"], + "created": created, + "accepted_at": ACCEPTED_AT, + } + for result in data["results"] + ], + } + + +def respond(method, path, *, body=None, **kwargs): + assert method == "POST" + if path.endswith("/upload/"): + return upload_response(body) + if path == BASE: + data = json.loads(body) + return experiment_response( + id=data["id"], + started_at=data["started_at"], + expected_item_count=data.get("expected_item_count"), + expected_result_count=data.get("expected_result_count"), + ) + action = path.rstrip("/").rsplit("/", 1)[1] + return experiment_response( + id=path.split("/")[-3], + created=False, + status={"complete": "completed", "fail": "failed"}[action], + finished_at=ACCEPTED_AT, + ) + + +@pytest.fixture(params=[False, True], ids=["sync", "async"]) +def client(request, monkeypatch): + transport = Mock(spec=SyncTransport) + if request.param: + transport = AsyncMock(spec=AsyncTransport) + monkeypatch.setattr(_client, "AsyncTransport", lambda *args: transport) + client = AsyncOfflineEvaluations(project_id=123, secret_key="test-secret") + else: + monkeypatch.setattr(_client, "SyncTransport", lambda *args: transport) + client = OfflineEvaluations(project_id=123, secret_key="test-secret") + transport.request.side_effect = respond + return client, transport + + +def invoke(method, *args, **kwargs): + result = method(*args, **kwargs) + return asyncio.run(result) if inspect.isawaitable(result) else result + + +def result(item=None, **kwargs): + return EvaluationResult( + item=item or EvaluationItem(input="question", output="answer"), + scorer_version_id=VERSION_ID, + **({"value": 1} | kwargs), + ) + + +def test_create_upload_and_explicit_completion(client): + evaluations, transport = client + experiment = invoke( + evaluations.create_experiment, + name="Quality benchmark", + run_source="ci", + expected_item_count=1, + expected_result_count=1, + application_version="commit-123", + ) + creation_body = json.loads(transport.request.call_args.kwargs["body"]) + assert str(UUID(experiment.id)) == creation_body["id"] + assert datetime.fromisoformat(creation_body["started_at"]).tzinfo is not None + assert creation_body["run_source"] == "ci" + assert creation_body["application_version"] == "commit-123" + assert experiment.receipt.expected_result_count == 1 + assert experiment.submission == creation_body + changed_copy = experiment.submission + changed_copy["name"] = "Changed locally" + assert experiment.submission == creation_body + + item = EvaluationItem(input="Question", output="Answer") + receipt = invoke( + experiment.upload_result, item=item, scorer_version_id=VERSION_ID, value=False + ) + assert receipt.items[0].id == item.id + assert receipt.results[0].created is True + assert receipt.results[0].accepted_at.tzinfo == timezone.utc + assert receipt.results[0].scorer_version_id == VERSION_ID + assert ( + json.loads(transport.request.call_args.kwargs["body"])["results"][0]["value"] + is False + ) + assert transport.request.call_args.args == ( + "POST", + BASE + experiment.id + "/upload/", + ) + completion = invoke(experiment.complete) + assert completion.status == "completed" + assert experiment.receipt == completion + assert transport.request.call_args.kwargs["body"] == b"{}" + + +def test_resume_is_local_and_supports_references_without_scorer_discovery(client): + evaluations, transport = client + experiment = evaluations.resume_experiment(UUID(EXPERIMENT_ID)) + assert experiment.id == EXPERIMENT_ID + assert experiment.submission is None + assert experiment.receipt is None + transport.request.assert_not_called() + item_id = str(uuid4()) + receipt = invoke( + experiment.upload_result, + item=item_id, + scorer_version_id=UUID(VERSION_ID), + value=0, + ) + body = json.loads(transport.request.call_args.kwargs["body"]) + assert body["items"] == [] + assert body["results"][0]["item_id"] == item_id + assert body["results"][0]["value"] == 0 + assert receipt.items[0].id == item_id + transport.request.assert_called_once() + + +def test_same_item_is_declared_once_for_multiple_scorers(client): + evaluations, transport = client + item = EvaluationItem(input={"question": "Paris?"}, output="France") + results = [ + result(item), + EvaluationResult(item=item, scorer_version_id=OTHER_VERSION_ID, value=True), + ] + receipt = invoke( + evaluations.resume_experiment(EXPERIMENT_ID).upload_results, results + ) + body = json.loads(transport.request.call_args.kwargs["body"]) + assert body["items"] == [item.to_dict()] + assert len(receipt.chunks) == 1 + assert len(receipt.chunks[0].results) == 2 + assert len(receipt.chunks[0].items) == 1 + + +def test_real_result_count_boundary_splits_without_dropping_declarations(client): + evaluations, transport = client + results = [result() for _ in range(1001)] + receipt = invoke( + evaluations.resume_experiment(EXPERIMENT_ID).upload_results, results + ) + assert [len(chunk.results) for chunk in receipt.chunks] == [1000, 1] + bodies = [ + json.loads(call.kwargs["body"]) for call in transport.request.call_args_list + ] + assert [len(body["items"]) for body in bodies] == [1000, 1] + assert [row["item_id"] for body in bodies for row in body["results"]] == [ + row.item_id for row in results + ] + + +def test_utf8_byte_boundary_splits_below_request_limit(client): + evaluations, transport = client + results = [result(EvaluationItem(input="🙂" * (190 * 1024))) for _ in range(7)] + receipt = invoke( + evaluations.resume_experiment(EXPERIMENT_ID).upload_results, results + ) + assert len(receipt.chunks) == 2 + assert sum(len(chunk.results) for chunk in receipt.chunks) == 7 + assert all( + len(call.kwargs["body"]) <= 5 * 1024 * 1024 + for call in transport.request.call_args_list + ) + + +def test_shared_item_is_redeclared_in_each_independently_replayable_chunk( + client, monkeypatch +): + evaluations, transport = client + monkeypatch.setattr(_client, "_MAX_RESULTS", 1) + item = EvaluationItem(input="Same execution") + results = [ + EvaluationResult(item=item.id, scorer_version_id=VERSION_ID, value=1), + EvaluationResult(item=item, scorer_version_id=OTHER_VERSION_ID, value=True), + ] + invoke(evaluations.resume_experiment(EXPERIMENT_ID).upload_results, results) + assert len(transport.request.call_args_list) == 2 + for call in transport.request.call_args_list: + assert json.loads(call.kwargs["body"])["items"] == [item.to_dict()] + + +@pytest.mark.parametrize("persistence", ["unknown", "rejected"]) +def test_partial_bulk_error_retains_acknowledged_chunks_and_exact_pending_work( + client, monkeypatch, persistence +): + evaluations, transport = client + monkeypatch.setattr(_client, "_MAX_RESULTS", 1) + results = [result() for _ in range(3)] + calls = 0 + + def fail_second(method, path, **kwargs): + nonlocal calls + calls += 1 + if calls == 2: + raise EvaluationAPIError( + status=503 if persistence == "unknown" else 409, persistence=persistence + ) + return respond(method, path, **kwargs) + + transport.request.side_effect = fail_second + experiment = evaluations.resume_experiment(EXPERIMENT_ID) + with pytest.raises(BulkUploadError) as caught: + invoke(experiment.upload_results, results) + error = caught.value + assert len(error.completed.chunks) == 1 + assert error.failed_index == 1 + assert error.pending_results == tuple(results[1:]) + assert error.persistence == persistence + failed_body = transport.request.call_args.kwargs["body"] + assert error.failed_chunk == json.loads(failed_body) + error.failed_chunk["results"].clear() + assert error.failed_chunk["results"] + assert transport.request.call_count == 2 + + transport.request.side_effect = respond + receipt = invoke(experiment.upload_results, error.pending_results) + assert len(receipt.chunks) == 2 + assert transport.request.call_args_list[2].kwargs["body"] == failed_body + assert not any( + call.args[1].endswith(("/fail/", "/complete/")) + for call in transport.request.call_args_list + ) + + +def test_creation_error_exposes_generated_identity_for_exact_retry(client): + evaluations, transport = client + transport.request.side_effect = EvaluationAPIError(persistence="unknown") + with pytest.raises(EvaluationAPIError) as caught: + invoke( + evaluations.create_experiment, + name="Recoverable", + run_source=None, + expected_result_count=0, + ) + original = transport.request.call_args.kwargs["body"] + saved = caught.value.submission + assert saved == json.loads(original) + assert saved["run_source"] is None + transport.request.side_effect = respond + experiment = invoke(evaluations.create_experiment, **saved) + assert experiment.id == saved["id"] + assert transport.request.call_args.kwargs["body"] == original + + +def test_malformed_creation_receipt_is_uncertain_and_recoverable(client): + evaluations, transport = client + transport.request.side_effect = None + transport.request.return_value = experiment_response(id=str(uuid4())) + with pytest.raises(EvaluationAPIError) as caught: + invoke(evaluations.create_experiment, name="Run", id=EXPERIMENT_ID) + assert caught.value.code == "invalid_response" + assert caught.value.persistence == "unknown" + assert caught.value.submission["id"] == EXPERIMENT_ID + + +@pytest.mark.parametrize( + "fault", + [ + "missing_items", + "missing_results", + "wrong_version", + "wrong_result_id", + "bad_created", + ], +) +def test_malformed_upload_receipts_cannot_report_delivery(client, fault): + evaluations, transport = client + + def malformed(method, path, *, body): + response = upload_response(body) + if fault == "missing_items": + response["items"] = [] + elif fault == "missing_results": + response["results"] = [] + elif fault == "wrong_version": + response["results"][0]["scorer_version_id"] = OTHER_VERSION_ID + elif fault == "wrong_result_id": + response["results"][0]["id"] = "invalid" + else: + response["results"][0]["created"] = "true" + return response + + transport.request.side_effect = malformed + with pytest.raises(BulkUploadError) as caught: + invoke(evaluations.resume_experiment(EXPERIMENT_ID).upload_results, [result()]) + assert caught.value.persistence == "unknown" + assert caught.value.completed.chunks == () + assert caught.value.failed_index == 0 + + +def test_duplicate_acknowledgments_and_server_count_conflicts_are_observable(client): + evaluations, transport = client + transport.request.side_effect = lambda method, path, *, body: upload_response( + body, created=False + ) + experiment = evaluations.resume_experiment(EXPERIMENT_ID) + receipt = invoke(experiment.upload_results, [result()]) + assert receipt.chunks[0].items[0].created is False + assert receipt.chunks[0].results[0].created is False + + error = EvaluationAPIError( + status=409, + code="expected_counts_mismatch", + response={"accepted_result_count": 1, "expected_result_count": 2}, + persistence="rejected", + ) + transport.request.side_effect = error + with pytest.raises(EvaluationAPIError) as caught: + invoke(experiment.complete) + assert caught.value is error + assert caught.value.response["expected_result_count"] == 2 + assert experiment.receipt is None + assert transport.request.call_count == 2 + + +def test_fail_is_explicit_and_uses_the_same_idempotent_lifecycle_contract(client): + evaluations, transport = client + receipt = invoke(evaluations.resume_experiment(EXPERIMENT_ID).fail) + assert receipt.status == "failed" + assert transport.request.call_args.args == ("POST", BASE + EXPERIMENT_ID + "/fail/") + assert transport.request.call_args.kwargs["body"] == b"{}" + + +def test_empty_bulk_and_local_validation_make_no_network_requests(client): + evaluations, transport = client + experiment = evaluations.resume_experiment(EXPERIMENT_ID) + assert invoke(experiment.upload_results, []).chunks == () + original = result() + with pytest.raises(ValueError, match="pair"): + invoke(experiment.upload_results, [original, original]) + changed = EvaluationItem(id=original.item_id, input="Changed execution") + conflict = EvaluationResult( + item=changed, scorer_version_id=OTHER_VERSION_ID, value=True + ) + with pytest.raises(ValueError, match="Conflicting"): + invoke(experiment.upload_results, [original, conflict]) + transport.request.assert_not_called() + + +def test_items_and_results_survive_json_restoration_without_mutable_aliases(): + content = {"messages": [{"text": "Original"}]} + item = EvaluationItem(input=content, output=None, metadata={"dataset": "v1"}) + categories = ["safe", "helpful"] + original = EvaluationResult( + item=item, + scorer_version_id=VERSION_ID, + value=categories, + reasoning="Two categories", + ) + saved = json.loads( + json.dumps({"item": item.to_dict(), "result": original.to_dict()}) + ) + content["messages"][0]["text"] = "Changed" + categories.append("other") + restored_item = EvaluationItem.from_dict(saved["item"]) + restored = EvaluationResult.from_dict(saved["result"], item=restored_item) + assert restored_item.to_dict() == item.to_dict() + assert restored.to_dict() == original.to_dict() + assert restored.item_id == item.id + assert restored.item is restored_item + item.to_dict()["payload"]["input"]["messages"].clear() + assert item.to_dict() == saved["item"] + with pytest.raises(ValueError, match="original id"): + EvaluationItem.from_dict({"payload": {}}) + with pytest.raises(ValueError, match="do not match"): + EvaluationResult.from_dict(saved["result"], item=EvaluationItem()) + + +def test_missing_empty_and_null_payloads_remain_distinct(): + omitted = EvaluationItem() + empty = EvaluationItem(payload={}) + null_fields = EvaluationItem(input=None, metadata=None) + assert "payload" not in omitted.to_dict() + assert empty.to_dict()["payload"] == {} + assert null_fields.to_dict()["payload"] == {"input": None, "metadata": None} + assert "payload" not in result().to_dict() + assert result(payload={}).to_dict()["payload"] == {} + assert result(reasoning=None).to_dict()["payload"] == {"reasoning": None} + with pytest.raises(ValueError, match="not both"): + EvaluationItem(input=None, payload={}) + + +@pytest.mark.parametrize("value", [False, True, 0, -1.25, ["safe", "helpful"]]) +def test_accepted_scores_preserve_type_and_polarity(value): + actual = result(value=value).to_dict()["value"] + if isinstance(value, bool): + assert actual is value + else: + assert actual == (sorted(value) if isinstance(value, list) else value) + + +@pytest.mark.parametrize("value", [2**53 - 1, 2**53, 2**53 + 2, 2**100]) +@pytest.mark.parametrize("sign", [1, -1]) +def test_exactly_representable_integer_scores_preserve_value(value, sign): + supplied = sign * value + actual = result(value=supplied).to_dict()["value"] + assert isinstance(actual, float) + assert actual == supplied + + +@pytest.mark.parametrize("value", [2**53 + 1, 2**100 + 1]) +@pytest.mark.parametrize("sign", [1, -1]) +def test_integer_scores_that_lose_precision_fail_before_upload(client, value, sign): + evaluations, transport = client + experiment = evaluations.resume_experiment(EXPERIMENT_ID) + with pytest.raises(ValueError, match="exactly representable"): + invoke( + experiment.upload_result, + item=EvaluationItem(), + scorer_version_id=VERSION_ID, + value=sign * value, + ) + transport.request.assert_not_called() + + +@pytest.mark.parametrize( + "value", + [ + None, + "pass", + [], + ["safe", "safe"], + [""], + ["a\x00b"], + float("nan"), + float("inf"), + 10**400, + ], +) +def test_invalid_scores_fail_before_upload(value): + with pytest.raises(ValueError): + result(value=value) + + +@pytest.mark.parametrize("status", ["error", "skipped", "not_applicable"]) +def test_non_ok_outcomes_have_no_score(status): + kwargs = {"status": status} + if status == "error": + kwargs |= {"error_code": "timeout", "error_message": "Judge timed out"} + outcome = EvaluationResult( + item=EvaluationItem(), scorer_version_id=VERSION_ID, **kwargs + ) + assert "value" not in outcome.to_dict() + with pytest.raises(ValueError, match="Only an ok"): + result(status=status, value=0) + + +@pytest.mark.parametrize( + "kwargs", + [ + {"error_code": "timeout"}, + {"error_message": "Timed out"}, + {"payload": {"error_message": "Timed out"}}, + ], +) +def test_error_fields_are_restricted_to_error_status(kwargs): + with pytest.raises(ValueError, match="Only an error"): + result(**kwargs) + + +def test_payload_byte_and_json_shape_validation_matches_backend_limits(): + with pytest.raises(ValueError, match="bytes"): + EvaluationItem(input="é" * (512 * 1024)) + with pytest.raises(ValueError, match="bytes"): + result(reasoning="x" * (256 * 1024)) + nested = {} + for _ in range(33): + nested = {"child": nested} + with pytest.raises(ValueError, match="nesting"): + EvaluationItem(input=nested) + for value in ( + {1: "not string"}, + {"a": object()}, + {"a": float("nan")}, + {"a": "\ud800"}, + ): + with pytest.raises(ValueError): + EvaluationItem(input=value) + + +def test_client_contexts_close_connections_without_closing_experiments(monkeypatch): + transport = Mock(spec=SyncTransport) + monkeypatch.setattr(_client, "SyncTransport", lambda *args: transport) + with OfflineEvaluations(project_id=123, secret_key="test") as evaluations: + evaluations.resume_experiment(EXPERIMENT_ID) + transport.close.assert_called_once() + transport.request.assert_not_called() + + async def run(): + transport = AsyncMock(spec=AsyncTransport) + monkeypatch.setattr(_client, "AsyncTransport", lambda *args: transport) + async with AsyncOfflineEvaluations( + project_id=123, secret_key="test" + ) as evaluations: + evaluations.resume_experiment(EXPERIMENT_ID) + transport.aclose.assert_awaited_once() + transport.request.assert_not_called() + + asyncio.run(run()) + + +def test_async_cancellation_never_completes_or_fails_experiment(monkeypatch): + async def run(): + transport = AsyncMock(spec=AsyncTransport) + transport.request.side_effect = asyncio.CancelledError + monkeypatch.setattr(_client, "AsyncTransport", lambda *args: transport) + async with AsyncOfflineEvaluations( + project_id=123, secret_key="test" + ) as evaluations: + experiment = evaluations.resume_experiment(EXPERIMENT_ID) + with pytest.raises(asyncio.CancelledError): + await experiment.upload_results([result()]) + assert experiment.receipt is None + transport.request.assert_awaited_once() + assert transport.request.call_args.args[1].endswith("/upload/") + transport.aclose.assert_awaited_once() + + asyncio.run(run()) + + +def test_example_persists_before_writes_and_replays_the_same_run(monkeypatch, tmp_path): + example = Path(__file__).parents[3] / "examples/example-offline-evaluations/main.py" + spec = importlib.util.spec_from_file_location("offline_evaluation_example", example) + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + state_path = tmp_path / "prepared-run.json" + monkeypatch.setenv("POSTHOG_PROJECT_ID", "123") + monkeypatch.setenv("POSTHOG_SECRET_KEY", "never-persist-this-secret") + monkeypatch.setenv("POSTHOG_HOST", "https://us.posthog.com") + monkeypatch.setenv("POSTHOG_SCORER_VERSION_ID", VERSION_ID) + monkeypatch.setattr("sys.argv", [str(example), "--state", str(state_path)]) + transport = Mock(spec=SyncTransport) + + def check_state_before_write(method, path, **kwargs): + assert state_path.exists() + assert "never-persist-this-secret" not in state_path.read_text() + return respond(method, path, **kwargs) + + transport.request.side_effect = check_state_before_write + monkeypatch.setattr(_client, "SyncTransport", lambda *args: transport) + module.main() + prepared_state = state_path.read_bytes() + module.main() + assert state_path.read_bytes() == prepared_state + requests = transport.request.call_args_list + assert len(requests) == 6 + assert requests[:3] == requests[3:] + assert transport.close.call_count == 2 diff --git a/posthog/test/ai/test_evaluations_scorers.py b/posthog/test/ai/test_evaluations_scorers.py new file mode 100644 index 000000000..128834c0f --- /dev/null +++ b/posthog/test/ai/test_evaluations_scorers.py @@ -0,0 +1,369 @@ +import asyncio +import copy +import inspect +import json +from datetime import timezone +from unittest.mock import AsyncMock, Mock +from uuid import UUID + +import pytest + +from posthog.ai.evaluations._errors import EvaluationAPIError +from posthog.ai.evaluations._scorer_types import ( + BooleanScorerConfig, + CategoricalPassingRule, + CategoricalScorerConfig, + CategoricalScorerOption, + NumericPassingRule, + NumericScorerConfig, +) +from posthog.ai.evaluations._scorers import AsyncScorers, Scorers +from posthog.ai.evaluations._transport import AsyncTransport, SyncTransport + + +SCORER_ID = "8cba6d9e-d2ee-4b23-a545-a61057da5870" +VERSION_ID = "f72481ec-1310-41de-b5bb-d20e1c59f121" +BASE = "/api/projects/123/llm_analytics/score_definitions/" + + +def scorer_response(**overrides): + return { + "id": SCORER_ID, + "name": "Quality", + "description": "", + "kind": "numeric", + "archived": False, + "current_version": 1, + "current_version_id": VERSION_ID, + "config": {"min": 0, "max": 1}, + "created_by": None, + "created_at": "2026-10-02T10:00:00Z", + "updated_at": "2026-10-02T10:00:00Z", + "team": 123, + **overrides, + } + + +def version_response(**overrides): + return { + "id": VERSION_ID, + "definition_id": SCORER_ID, + "version": 1, + "kind": "numeric", + "config": {"min": 0, "max": 1}, + "created_at": "2026-10-02T10:00:00Z", + "created_by": None, + **overrides, + } + + +@pytest.fixture(params=[False, True], ids=["sync", "async"]) +def scorers(request): + transport = Mock(spec=SyncTransport) + namespace = Scorers(transport, 123) + if request.param: + transport = AsyncMock(spec=AsyncTransport) + namespace = AsyncScorers(transport, 123) + transport.request.return_value = scorer_response() + return namespace, transport + + +def invoke(method, *args, **kwargs): + result = method(*args, **kwargs) + return asyncio.run(result) if inspect.isawaitable(result) else result + + +@pytest.mark.parametrize( + ("kind", "config"), + [ + ( + "boolean", + BooleanScorerConfig( + true_is_failure=True, true_label="Unsafe", false_label="Safe" + ), + ), + ("boolean", BooleanScorerConfig(true_is_failure=False)), + ("boolean", BooleanScorerConfig(true_is_failure=None)), + ( + "numeric", + NumericScorerConfig( + min=0, + max=1, + step=0.01, + passing_rule=NumericPassingRule(operator="gte", threshold=0.8), + ), + ), + ( + "numeric", + NumericScorerConfig( + min=None, + max=None, + step=None, + passing_rule=NumericPassingRule(operator="lte", threshold=0), + ), + ), + ("numeric", NumericScorerConfig(passing_rule=None)), + ( + "categorical", + CategoricalScorerConfig( + options=[ + CategoricalScorerOption(key="safe", label="Safe"), + CategoricalScorerOption(key="helpful", label="Helpful"), + ], + selection_mode="multiple", + min_selections=1, + max_selections=2, + passing_rule=CategoricalPassingRule(categories=["safe", "helpful"]), + ), + ), + ( + "categorical", + CategoricalScorerConfig( + options=[CategoricalScorerOption(key="safe", label="Safe")], + selection_mode="single", + min_selections=None, + max_selections=None, + passing_rule=None, + ), + ), + ], +) +def test_typed_configs_preserve_polarity_and_nullable_fields(scorers, kind, config): + namespace, transport = scorers + transport.request.return_value = scorer_response(kind=kind, config=config) + scorer = invoke(namespace.create, name="Quality", kind=kind, config=config) + + assert json.loads(transport.request.call_args.kwargs["body"]) == { + "name": "Quality", + "kind": kind, + "config": config, + } + assert scorer.current_version_id == VERSION_ID + assert scorer.config == config + assert scorer.created_at.tzinfo == timezone.utc + assert transport.request.call_args.args == ("POST", BASE) + assert transport.request.call_args.kwargs["retry_safe"] is False + + +@pytest.mark.parametrize("description", [None, "", "説明"]) +def test_create_preserves_explicit_description(scorers, description): + namespace, transport = scorers + invoke( + namespace.create, + name="Quality", + kind="numeric", + config=NumericScorerConfig(), + description=description, + ) + assert ( + json.loads(transport.request.call_args.kwargs["body"])["description"] + == description + ) + + +def test_metadata_update_preserves_omission_null_and_false(scorers): + namespace, transport = scorers + invoke(namespace.update, SCORER_ID, description=None, archived=False) + assert transport.request.call_args.args == ("PATCH", BASE + SCORER_ID + "/") + assert json.loads(transport.request.call_args.kwargs["body"]) == { + "description": None, + "archived": False, + } + assert transport.request.call_args.kwargs["retry_safe"] is False + invoke(namespace.update, SCORER_ID, name="Renamed") + assert json.loads(transport.request.call_args.kwargs["body"]) == {"name": "Renamed"} + + +def test_version_bump_preserves_configuration_and_concurrency_guard(scorers): + namespace, transport = scorers + config = NumericScorerConfig(min=0, max=1) + transport.request.return_value = scorer_response(current_version=2) + updated = invoke( + namespace.create_version, + UUID(SCORER_ID), + config=config, + base_version=1, + name="New name", + description=None, + ) + assert updated.current_version == 2 + assert updated.current_version_id == VERSION_ID + assert transport.request.call_args.args == ( + "POST", + BASE + SCORER_ID + "/new_version/", + ) + assert json.loads(transport.request.call_args.kwargs["body"]) == { + "config": config, + "base_version": 1, + "name": "New name", + "description": None, + } + assert transport.request.call_args.kwargs["retry_safe"] is False + + +def test_stale_version_is_not_refreshed_or_retried(scorers): + namespace, transport = scorers + error = EvaluationAPIError( + status=409, response={"current_version": 5}, persistence="rejected" + ) + transport.request.side_effect = error + with pytest.raises(EvaluationAPIError) as caught: + invoke( + namespace.create_version, + SCORER_ID, + config=NumericScorerConfig(), + base_version=4, + ) + assert caught.value is error + assert caught.value.response == {"current_version": 5} + transport.request.assert_called_once() + + +def test_scorer_list_supports_filters_and_explicit_offset_continuation(scorers): + namespace, transport = scorers + transport.request.return_value = { + "count": 3, + "next": "https://us.posthog.com" + BASE + "?limit=1&offset=2&kind=numeric", + "previous": "https://us.posthog.com" + BASE + "?limit=1&kind=numeric", + "results": [scorer_response()], + } + page = invoke( + namespace.list, + limit=1, + offset=1, + search="Quality", + kind="numeric", + archived=False, + order_by="-current_version", + ) + assert page.count == 3 + assert page.next_offset == 2 + assert page.results[0].id == SCORER_ID + assert transport.request.call_args.kwargs["params"] == { + "limit": 1, + "offset": 1, + "search": "Quality", + "kind": "numeric", + "archived": "false", + "order_by": "-current_version", + } + assert transport.request.call_args.args == ("GET", BASE) + transport.request.return_value = { + "count": 0, + "next": None, + "previous": None, + "results": [], + } + assert invoke(namespace.list).next_offset is None + assert "archived" not in transport.request.call_args.kwargs["params"] + + +def test_version_history_uses_cursors_instead_of_offsets(scorers): + namespace, transport = scorers + transport.request.return_value = { + "count": 2, + "next_cursor": "opaque==", + "results": [version_response()], + } + page = invoke(namespace.list_versions, SCORER_ID, limit=1, cursor="prior==") + assert page.next_cursor == "opaque==" + assert page.results[0].id == VERSION_ID + assert page.results[0].definition_id == SCORER_ID + assert transport.request.call_args.args == ("GET", BASE + SCORER_ID + "/versions/") + assert transport.request.call_args.kwargs["params"] == { + "limit": 1, + "cursor": "prior==", + } + transport.request.return_value = version_response() + version = invoke(namespace.get_version, UUID(SCORER_ID), UUID(VERSION_ID)) + assert version.id == VERSION_ID + assert transport.request.call_args.args == ( + "GET", + BASE + SCORER_ID + "/versions/" + VERSION_ID + "/", + ) + + +@pytest.mark.parametrize( + "invalid_id", ["Quality", "../other-project", "?limit=10", "", 123] +) +def test_definition_paths_require_uuid_and_never_lookup_by_name(scorers, invalid_id): + namespace, transport = scorers + with pytest.raises(ValueError, match="UUID"): + invoke(namespace.get, invalid_id) + with pytest.raises(ValueError, match="UUID"): + invoke(namespace.get_version, SCORER_ID, invalid_id) + transport.request.assert_not_called() + + +@pytest.mark.parametrize("invalid_limit", [0, -1, 101, True]) +def test_version_page_limit_rejected_before_io(scorers, invalid_limit): + namespace, transport = scorers + with pytest.raises(ValueError): + invoke(namespace.list_versions, SCORER_ID, limit=invalid_limit) + transport.request.assert_not_called() + + +@pytest.mark.parametrize( + "override", + [ + {"id": "invalid"}, + {"archived": "false"}, + {"created_at": "2026-10-02T12:00:00"}, + {"config": []}, + {"current_version": True}, + {"current_version_id": None}, + ], +) +def test_malformed_success_preserves_unknown_persistence(scorers, override): + namespace, transport = scorers + transport.request.return_value = scorer_response(**override) + with pytest.raises(EvaluationAPIError) as caught: + invoke( + namespace.create, + name="Quality", + kind="numeric", + config=NumericScorerConfig(), + ) + assert caught.value.code == "invalid_response" + assert caught.value.persistence == "unknown" + assert "Quality" not in str(caught.value) + transport.request.assert_called_once() + + +def test_configuration_response_is_detached_from_transport_and_creation_input(scorers): + namespace, transport = scorers + config = CategoricalScorerConfig( + options=[CategoricalScorerOption(key="a", label="A")] + ) + response = scorer_response(kind="categorical", config=config) + transport.request.return_value = response + scorer = invoke(namespace.create, name="Quality", kind="categorical", config=config) + expected = copy.deepcopy(config) + config["options"][0]["label"] = "Changed" + assert scorer.config == expected + assert json.loads(transport.request.call_args.kwargs["body"])["config"] == expected + + +def test_nonfinite_configuration_is_never_sent(scorers): + namespace, transport = scorers + with pytest.raises(ValueError): + invoke( + namespace.create, + name="Quality", + kind="numeric", + config=NumericScorerConfig(min=float("nan")), + ) + transport.request.assert_not_called() + + +def test_async_cancellation_propagates_without_followup_writes(): + async def run(): + transport = AsyncMock(spec=AsyncTransport) + transport.request.side_effect = asyncio.CancelledError + with pytest.raises(asyncio.CancelledError): + await AsyncScorers(transport, 123).create( + name="Quality", kind="numeric", config=NumericScorerConfig() + ) + transport.request.assert_called_once() + + asyncio.run(run()) diff --git a/posthog/test/ai/test_evaluations_transport.py b/posthog/test/ai/test_evaluations_transport.py new file mode 100644 index 000000000..3f62fe932 --- /dev/null +++ b/posthog/test/ai/test_evaluations_transport.py @@ -0,0 +1,387 @@ +from __future__ import annotations + +import asyncio +import json +import subprocess +import sys +from datetime import datetime, timedelta, timezone +from email.utils import format_datetime +from unittest.mock import AsyncMock, patch + +import httpx +import pytest +import requests + +from posthog.ai.evaluations._errors import EvaluationAPIError +from posthog.ai.evaluations._transport import AsyncTransport, SyncTransport + +PATH = "/api/projects/123/ai_observability/offline_experiments/" +BODY = b'{"id":"stable-id","name":"evaluation data"}' + + +def response(status=200, payload=None, headers=None): + result = requests.Response() + result.status_code = status + result._content = json.dumps(payload if payload is not None else {}).encode() + result.headers.update(headers or {}) + return result + + +@pytest.fixture +def transport(): + client = SyncTransport(123, "phx_test_secret") + yield client + client.close() + + +@pytest.mark.parametrize( + "host,expected", + [ + (None, "https://us.posthog.com"), + ("https://app.posthog.com/", "https://us.posthog.com"), + ("https://us.i.posthog.com", "https://us.posthog.com"), + ("https://eu.i.posthog.com/", "https://eu.posthog.com"), + ("https://eu.posthog.com", "https://eu.posthog.com"), + ("http://localhost:8000/posthog/", "http://localhost:8000/posthog"), + ("https://proxy.example.com/prefix", "https://proxy.example.com/prefix"), + ], +) +def test_management_hosts_and_authenticated_request(host, expected): + transport = SyncTransport(123, "phx_test_secret", host=host, timeout=7) + with patch.object(transport._session, "send", return_value=response()) as send: + assert transport.request("POST", PATH, body=BODY) == {} + sent = send.call_args.args[0] + assert sent.url == expected + PATH + assert sent.headers["Authorization"] == "Bearer phx_test_secret" + assert sent.headers["Content-Type"] == "application/json" + assert sent.headers["User-Agent"].startswith("posthog-python/") + assert sent.body == BODY + assert send.call_args.kwargs["allow_redirects"] is False + assert send.call_args.kwargs["timeout"] == 7 + assert transport._session.adapters["https://"].max_retries.total == 0 + transport.close() + + +@pytest.mark.parametrize( + "options", + [ + {"project_id": 0}, + {"project_id": True}, + {"secret_key": ""}, + {"secret_key": "with\nnewline"}, + {"host": "example.com"}, + {"host": "https://user:password@example.com"}, + {"host": "https://example.com?key=secret"}, + {"timeout": 0}, + {"timeout": float("inf")}, + {"max_retries": -1}, + {"max_retries": True}, + ], +) +@pytest.mark.parametrize("transport_class", [SyncTransport, AsyncTransport]) +def test_invalid_configuration_fails_before_io(options, transport_class): + with pytest.raises(ValueError): + transport_class(**{"project_id": 123, "secret_key": "secret", **options}) + + +def test_lost_response_retries_identical_body_and_identity(transport): + with ( + patch.object( + transport._session, + "request", + side_effect=[ + requests.ReadTimeout("secret message"), + response(201, {"id": "stable-id"}), + ], + ) as request, + patch("posthog.ai.evaluations._transport.time.sleep") as sleep, + ): + assert transport.request("POST", PATH, body=BODY) == {"id": "stable-id"} + assert request.call_count == 2 + assert request.call_args_list[0] == request.call_args_list[1] + sleep.assert_called_once() + + +@pytest.mark.parametrize("failure", [requests.ReadTimeout(), response(503)]) +def test_non_idempotent_mutation_never_replays(transport, failure): + with ( + patch.object(transport._session, "request", side_effect=[failure]) as request, + patch("posthog.ai.evaluations._transport.time.sleep") as sleep, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY, retry_safe=False) + assert request.call_count == 1 + assert caught.value.persistence == "unknown" + sleep.assert_not_called() + + +def test_exhausted_transport_retry_budget_is_unknown_and_payload_free(transport): + with ( + patch.object( + transport._session, + "request", + side_effect=requests.ConnectionError("phx_test_secret evaluation data"), + ) as request, + patch("posthog.ai.evaluations._transport.time.sleep") as sleep, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY) + assert request.call_count == 4 + assert sleep.call_count == 3 + assert caught.value.status is None + assert caught.value.code == "transport_error" + assert caught.value.persistence == "unknown" + assert "phx_test_secret" not in repr(caught.value) + assert "evaluation data" not in str(caught.value) + assert caught.value.__cause__ is None + + +def test_errors_retain_server_fields_and_conflict_counts_without_logging_them( + transport, +): + payload = { + "code": "incomplete", + "detail": "evaluation data and secret", + "attr": "expected_result_count", + "expected_result_count": 3, + "result_count": 2, + "errors": [{"attr": "items.0.payload", "detail": "private input"}], + } + with ( + patch.object( + transport._session, "request", return_value=response(409, payload) + ) as request, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY) + error = caught.value + assert request.call_count == 1 + assert error.status == 409 + assert error.code == "incomplete" + assert error.detail == payload["detail"] + assert error.attr == "expected_result_count" + assert error.errors == payload["errors"] + assert error.response == payload + assert error.persistence == "rejected" + assert "evaluation data" not in str(error) + assert "private input" not in repr(error) + + +def test_rejection_after_lost_response_does_not_claim_definite_rejection(transport): + with ( + patch.object( + transport._session, + "request", + side_effect=[requests.ReadTimeout(), response(409, {"code": "conflict"})], + ), + patch("posthog.ai.evaluations._transport.time.sleep"), + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY) + assert caught.value.status == 409 + assert caught.value.persistence == "unknown" + + +@pytest.mark.parametrize("status", [200, 201]) +@pytest.mark.parametrize("content", [b"not JSON", b"[]", b"null"]) +def test_malformed_success_is_unknown_and_not_retried(transport, status, content): + malformed = response(status) + malformed._content = content + with ( + patch.object(transport._session, "request", return_value=malformed) as request, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY) + assert request.call_count == 1 + assert caught.value.code == "invalid_response" + assert caught.value.persistence == "unknown" + + +@pytest.mark.parametrize("retry_after_format", ["seconds", "date"]) +def test_retry_after_is_minimum_delay(transport, retry_after_format): + retry_after = ( + "2" + if retry_after_format == "seconds" + else format_datetime(datetime.now(timezone.utc) + timedelta(seconds=20)) + ) + with ( + patch.object( + transport._session, + "request", + side_effect=[ + response(429, headers={"Retry-After": retry_after}), + response(), + ], + ), + patch("posthog.ai.evaluations._transport.time.sleep") as sleep, + ): + transport.request("POST", PATH, body=BODY) + sleep.assert_called_once() + assert 2 <= sleep.call_args.args[0] <= 30 + + +def test_retry_after_longer_than_budget_surfaces_without_retrying_early(transport): + with ( + patch.object( + transport._session, + "request", + return_value=response(429, headers={"Retry-After": "3600"}), + ) as request, + patch("posthog.ai.evaluations._transport.time.sleep") as sleep, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY) + assert request.call_count == 1 + assert caught.value.retry_after == 3600 + assert caught.value.persistence == "rejected" + sleep.assert_not_called() + + +def test_redirect_is_never_followed(transport): + with ( + patch.object( + transport._session, + "request", + return_value=response( + 307, headers={"Location": "https://elsewhere.example"} + ), + ) as request, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("POST", PATH, body=BODY) + assert request.call_count == 1 + assert request.call_args.kwargs["allow_redirects"] is False + assert caught.value.status == 307 + + +def test_closed_client_does_not_send(transport): + transport.close() + with ( + patch.object(transport._session, "request") as request, + pytest.raises(EvaluationAPIError) as caught, + ): + transport.request("GET", PATH) + assert caught.value.persistence == "not_sent" + request.assert_not_called() + + +@pytest.mark.asyncio +async def test_async_native_request_uses_auth_host_and_preserves_bytes(): + seen = [] + + async def handle(request): + seen.append(request) + if len(seen) == 1: + raise httpx.ReadTimeout("phx_test_secret") + return httpx.Response(201, json={"id": "stable-id"}) + + real_client = httpx.AsyncClient + with patch( + "httpx.AsyncClient", + side_effect=lambda **kwargs: real_client( + transport=httpx.MockTransport(handle), **kwargs + ), + ) as factory: + transport = AsyncTransport( + 123, "phx_test_secret", host="https://eu.i.posthog.com", timeout=9 + ) + with patch( + "posthog.ai.evaluations._transport.asyncio.sleep", new_callable=AsyncMock + ) as sleep: + assert await transport.request("POST", PATH, body=BODY) == {"id": "stable-id"} + assert len(seen) == 2 + assert seen[0].content == seen[1].content == BODY + assert str(seen[0].url) == "https://eu.posthog.com" + PATH + assert seen[0].headers["Authorization"] == "Bearer phx_test_secret" + assert factory.call_args.kwargs["follow_redirects"] is False + assert factory.call_args.kwargs["timeout"] == 9 + sleep.assert_awaited_once() + await transport.aclose() + assert transport._client.is_closed + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "failure", + [httpx.ReadTimeout("secret"), httpx.Response(503), asyncio.CancelledError()], +) +async def test_async_non_idempotent_failure_and_cancellation_never_replay(failure): + transport = AsyncTransport(123, "phx_test_secret") + with ( + patch.object(transport._client, "request", side_effect=[failure]) as request, + patch( + "posthog.ai.evaluations._transport.asyncio.sleep", new_callable=AsyncMock + ) as sleep, + pytest.raises( + asyncio.CancelledError + if isinstance(failure, asyncio.CancelledError) + else EvaluationAPIError + ), + ): + await transport.request("POST", PATH, body=BODY, retry_safe=False) + request.assert_awaited_once() + sleep.assert_not_awaited() + await transport.aclose() + + +@pytest.mark.asyncio +async def test_async_rate_limit_larger_than_wait_budget_surfaces(): + transport = AsyncTransport(123, "phx_test_secret") + with ( + patch.object( + transport._client, + "request", + return_value=httpx.Response(429, headers={"Retry-After": "60"}), + ) as request, + patch( + "posthog.ai.evaluations._transport.asyncio.sleep", new_callable=AsyncMock + ) as sleep, + pytest.raises(EvaluationAPIError) as caught, + ): + await transport.request("POST", PATH, body=BODY) + request.assert_awaited_once() + sleep.assert_not_awaited() + assert caught.value.status == 429 + assert caught.value.retry_after == 60 + await transport.aclose() + + +@pytest.mark.asyncio +async def test_async_closed_client_does_not_send(): + transport = AsyncTransport(123, "phx_test_secret") + await transport.aclose() + with ( + patch.object(transport._client, "request") as request, + pytest.raises(EvaluationAPIError) as caught, + ): + await transport.request("GET", PATH) + assert caught.value.persistence == "not_sent" + request.assert_not_awaited() + + +def test_optional_httpx_is_required_only_for_async_construction(): + script = """ +import sys +sys.modules['httpx'] = None +from posthog.ai.evaluations._transport import SyncTransport, AsyncTransport +client = SyncTransport(123, 'secret') +client.close() +try: + AsyncTransport(123, 'secret') +except RuntimeError as error: + assert 'posthog[async]' in str(error) +else: + raise AssertionError('Expected optional dependency guidance') +""" + subprocess.run([sys.executable, "-c", script], check=True, timeout=15) + + +def test_error_snapshots_response_and_submission(): + response_data = {"errors": [{"attr": "input", "detail": "private input"}]} + submission = {"id": "stable-id", "payload": {"input": "private input"}} + error = EvaluationAPIError(response=response_data, submission=submission) + response_data["errors"].clear() + submission["payload"]["input"] = "changed" + assert error.response["errors"] == [{"attr": "input", "detail": "private input"}] + assert error.submission["payload"]["input"] == "private input" + assert "private input" not in str(error) diff --git a/pyproject.toml b/pyproject.toml index e06ef657d..28df07926 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -105,6 +105,7 @@ test = [ packages = [ "posthog", "posthog.ai", + "posthog.ai.evaluations", "posthog.ai.langchain", "posthog.ai.openai", "posthog.ai.openai_agents", diff --git a/references/public_api_snapshot.txt b/references/public_api_snapshot.txt index efcc1f741..017fbb380 100644 --- a/references/public_api_snapshot.txt +++ b/references/public_api_snapshot.txt @@ -386,6 +386,122 @@ attribute posthog.ai.anthropic.anthropic_providers.AsyncAnthropicBedrock.message attribute posthog.ai.anthropic.anthropic_providers.AsyncAnthropicVertex.messages = AsyncWrappedMessages(self) attribute posthog.ai.claude_agent_sdk.client.log = logging.getLogger('posthog') attribute posthog.ai.claude_agent_sdk.processor.log = logging.getLogger('posthog') +attribute posthog.ai.evaluations.AsyncExperiment.id: str +attribute posthog.ai.evaluations.AsyncExperiment.receipt = receipt +attribute posthog.ai.evaluations.AsyncExperiment.submission: dict[str, Any] | None +attribute posthog.ai.evaluations.AsyncOfflineEvaluations.scorers = AsyncScorers(self._transport, project_id) +attribute posthog.ai.evaluations.BooleanScorerConfig.false_label: str +attribute posthog.ai.evaluations.BooleanScorerConfig.true_is_failure: bool | None +attribute posthog.ai.evaluations.BooleanScorerConfig.true_label: str +attribute posthog.ai.evaluations.BulkUploadError.completed = completed +attribute posthog.ai.evaluations.BulkUploadError.failed_chunk: dict[str, Any] +attribute posthog.ai.evaluations.BulkUploadError.failed_index = failed_index +attribute posthog.ai.evaluations.BulkUploadError.pending_results = pending_results +attribute posthog.ai.evaluations.BulkUploadReceipt.chunks: tuple[UploadReceipt, ...] +attribute posthog.ai.evaluations.CategoricalPassingRule.categories: list[str] +attribute posthog.ai.evaluations.CategoricalScorerConfig.max_selections: int | None +attribute posthog.ai.evaluations.CategoricalScorerConfig.min_selections: int | None +attribute posthog.ai.evaluations.CategoricalScorerConfig.options: Required[list[CategoricalScorerOption]] +attribute posthog.ai.evaluations.CategoricalScorerConfig.passing_rule: CategoricalPassingRule | None +attribute posthog.ai.evaluations.CategoricalScorerConfig.selection_mode: Literal['single', 'multiple'] +attribute posthog.ai.evaluations.CategoricalScorerOption.key: str +attribute posthog.ai.evaluations.CategoricalScorerOption.label: str +attribute posthog.ai.evaluations.EvaluationAPIError.attr = attr +attribute posthog.ai.evaluations.EvaluationAPIError.code = code +attribute posthog.ai.evaluations.EvaluationAPIError.detail = deepcopy(detail) +attribute posthog.ai.evaluations.EvaluationAPIError.errors = deepcopy(errors) +attribute posthog.ai.evaluations.EvaluationAPIError.persistence = persistence +attribute posthog.ai.evaluations.EvaluationAPIError.response = deepcopy(response) +attribute posthog.ai.evaluations.EvaluationAPIError.retry_after = retry_after +attribute posthog.ai.evaluations.EvaluationAPIError.status = status +attribute posthog.ai.evaluations.EvaluationAPIError.submission = deepcopy(submission) +attribute posthog.ai.evaluations.EvaluationItem.id: str +attribute posthog.ai.evaluations.EvaluationResult.item: EvaluationItem | str +attribute posthog.ai.evaluations.EvaluationResult.item_id: str +attribute posthog.ai.evaluations.EvaluationResult.scorer_version_id: str +attribute posthog.ai.evaluations.Experiment.id: str +attribute posthog.ai.evaluations.Experiment.receipt = receipt +attribute posthog.ai.evaluations.Experiment.submission: dict[str, Any] | None +attribute posthog.ai.evaluations.ExperimentOptions.application_version: str | None +attribute posthog.ai.evaluations.ExperimentOptions.dataset_identifier: str | None +attribute posthog.ai.evaluations.ExperimentOptions.dataset_revision_id: str | UUID | None +attribute posthog.ai.evaluations.ExperimentOptions.dataset_revision_identifier: str | None +attribute posthog.ai.evaluations.ExperimentOptions.dataset_source: str | None +attribute posthog.ai.evaluations.ExperimentOptions.expected_item_count: int | None +attribute posthog.ai.evaluations.ExperimentOptions.expected_result_count: int | None +attribute posthog.ai.evaluations.ExperimentOptions.id: str | UUID +attribute posthog.ai.evaluations.ExperimentOptions.model_version: str | None +attribute posthog.ai.evaluations.ExperimentOptions.prompt_version: str | None +attribute posthog.ai.evaluations.ExperimentOptions.run_source: RunSource | None +attribute posthog.ai.evaluations.ExperimentOptions.started_at: datetime | str +attribute posthog.ai.evaluations.ExperimentOptions.suite_key: str | None +attribute posthog.ai.evaluations.ExperimentReceipt.accepted_item_count: int +attribute posthog.ai.evaluations.ExperimentReceipt.accepted_result_count: int +attribute posthog.ai.evaluations.ExperimentReceipt.created: bool +attribute posthog.ai.evaluations.ExperimentReceipt.created_at: datetime +attribute posthog.ai.evaluations.ExperimentReceipt.expected_item_count: int | None +attribute posthog.ai.evaluations.ExperimentReceipt.expected_result_count: int | None +attribute posthog.ai.evaluations.ExperimentReceipt.finished_at: datetime | None +attribute posthog.ai.evaluations.ExperimentReceipt.id: str +attribute posthog.ai.evaluations.ExperimentReceipt.started_at: datetime +attribute posthog.ai.evaluations.ExperimentReceipt.status: Literal['uploading', 'completed', 'failed'] +attribute posthog.ai.evaluations.ItemReceipt.accepted_at: datetime +attribute posthog.ai.evaluations.ItemReceipt.created: bool +attribute posthog.ai.evaluations.ItemReceipt.id: str +attribute posthog.ai.evaluations.JSONValue = Union[None, bool, int, float, str, list['JSONValue'], dict[str, 'JSONValue']] +attribute posthog.ai.evaluations.NumericPassingRule.operator: Literal['gte', 'lte'] +attribute posthog.ai.evaluations.NumericPassingRule.threshold: float +attribute posthog.ai.evaluations.NumericScorerConfig.max: float | None +attribute posthog.ai.evaluations.NumericScorerConfig.min: float | None +attribute posthog.ai.evaluations.NumericScorerConfig.passing_rule: NumericPassingRule | None +attribute posthog.ai.evaluations.NumericScorerConfig.step: float | None +attribute posthog.ai.evaluations.OfflineEvaluations.scorers = Scorers(self._transport, project_id) +attribute posthog.ai.evaluations.ResultOptions.error_code: str | None +attribute posthog.ai.evaluations.ResultOptions.error_message: str | None +attribute posthog.ai.evaluations.ResultOptions.evaluated_at: datetime | str | None +attribute posthog.ai.evaluations.ResultOptions.evaluator_trace_id: str | None +attribute posthog.ai.evaluations.ResultOptions.metadata: dict[str, JSONValue] | None +attribute posthog.ai.evaluations.ResultOptions.payload: dict[str, JSONValue] +attribute posthog.ai.evaluations.ResultOptions.reasoning: str | None +attribute posthog.ai.evaluations.ResultOptions.status: ResultStatus +attribute posthog.ai.evaluations.ResultOptions.value: ScoreValue | None +attribute posthog.ai.evaluations.ResultReceipt.item_id: str +attribute posthog.ai.evaluations.ResultReceipt.scorer_version_id: str +attribute posthog.ai.evaluations.ResultStatus = Literal['ok', 'error', 'skipped', 'not_applicable'] +attribute posthog.ai.evaluations.RunSource = Literal['ci', 'local', 'scheduled'] +attribute posthog.ai.evaluations.ScoreValue = bool | int | float | list[str] +attribute posthog.ai.evaluations.Scorer.archived: bool +attribute posthog.ai.evaluations.Scorer.config: _ConfigT +attribute posthog.ai.evaluations.Scorer.created_at: datetime +attribute posthog.ai.evaluations.Scorer.created_by: dict[str, Any] | None +attribute posthog.ai.evaluations.Scorer.current_version: int +attribute posthog.ai.evaluations.Scorer.current_version_id: str +attribute posthog.ai.evaluations.Scorer.description: str +attribute posthog.ai.evaluations.Scorer.id: str +attribute posthog.ai.evaluations.Scorer.kind: ScorerKind +attribute posthog.ai.evaluations.Scorer.name: str +attribute posthog.ai.evaluations.Scorer.team: int +attribute posthog.ai.evaluations.Scorer.updated_at: datetime +attribute posthog.ai.evaluations.ScorerConfig = BooleanScorerConfig | NumericScorerConfig | CategoricalScorerConfig +attribute posthog.ai.evaluations.ScorerKind = Literal['boolean', 'numeric', 'categorical'] +attribute posthog.ai.evaluations.ScorerOrder = Literal['name', '-name', 'kind', '-kind', 'created_at', '-created_at', 'updated_at', '-updated_at', 'current_version', '-current_version'] +attribute posthog.ai.evaluations.ScorerPage.count: int +attribute posthog.ai.evaluations.ScorerPage.next: str | None +attribute posthog.ai.evaluations.ScorerPage.next_offset: int | None +attribute posthog.ai.evaluations.ScorerPage.previous: str | None +attribute posthog.ai.evaluations.ScorerPage.results: tuple[Scorer[ScorerConfig], ...] +attribute posthog.ai.evaluations.ScorerVersion.config: _ConfigT +attribute posthog.ai.evaluations.ScorerVersion.created_at: datetime +attribute posthog.ai.evaluations.ScorerVersion.created_by: dict[str, Any] | None +attribute posthog.ai.evaluations.ScorerVersion.definition_id: str +attribute posthog.ai.evaluations.ScorerVersion.id: str +attribute posthog.ai.evaluations.ScorerVersion.kind: ScorerKind +attribute posthog.ai.evaluations.ScorerVersion.version: int +attribute posthog.ai.evaluations.ScorerVersionPage.count: int +attribute posthog.ai.evaluations.ScorerVersionPage.next_cursor: str | None +attribute posthog.ai.evaluations.ScorerVersionPage.results: tuple[ScorerVersion[ScorerConfig], ...] +attribute posthog.ai.evaluations.UploadReceipt.items: tuple[ItemReceipt, ...] +attribute posthog.ai.evaluations.UploadReceipt.results: tuple[ResultReceipt, ...] attribute posthog.ai.gateway.POSTHOG_AI_GATEWAY_HOSTS = ['gateway.posthog.com', 'gateway.us.posthog.com', 'gateway.eu.posthog.com', 'ai-gateway.us.posthog.com', 'ai-gateway.eu.posthog.com'] attribute posthog.ai.gateway.log = logging.getLogger('posthog') attribute posthog.ai.gemini.gemini.Client.aio = _GeminiAioNamespace(models=(AsyncModels(provider_client=(self._provider_client), posthog_client=(self._ph_client), posthog_distinct_id=posthog_distinct_id, posthog_properties=posthog_properties, posthog_privacy_mode=posthog_privacy_mode, posthog_groups=posthog_groups, **kwargs)), files=(self._provider_client.aio.files)) @@ -977,6 +1093,33 @@ class posthog.ai.anthropic.anthropic_providers.AsyncAnthropicBedrock(posthog_cli class posthog.ai.anthropic.anthropic_providers.AsyncAnthropicVertex(posthog_client: Optional[PostHogClient] = None, **kwargs) class posthog.ai.claude_agent_sdk.client.PostHogClaudeSDKClient(options: Optional[ClaudeAgentOptions] = None, transport: Any = None, *, posthog_client: Optional[Client] = None, posthog_distinct_id: Optional[Union[str, Callable[[ResultMessage], Optional[str]]]] = None, posthog_trace_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: bool = False, posthog_groups: Optional[Dict[str, Any]] = None) class posthog.ai.claude_agent_sdk.processor.PostHogClaudeAgentProcessor(client: Optional[Client] = None, distinct_id: Optional[Union[str, Callable[[ResultMessage], Optional[str]]]] = None, privacy_mode: bool = False, groups: Optional[Dict[str, Any]] = None, properties: Optional[Dict[str, Any]] = None) +class posthog.ai.evaluations.AsyncExperiment(transport: AsyncTransport, project_id: int, id: str, *, receipt: ExperimentReceipt | None = None, submission: bytes | None = None) +class posthog.ai.evaluations.AsyncOfflineEvaluations(*, project_id: int, secret_key: str, host: str | None = None, timeout: float = 15, max_retries: int = 3) +class posthog.ai.evaluations.AsyncScorers(transport: AsyncTransport, project_id: int) +class posthog.ai.evaluations.BooleanScorerConfig +class posthog.ai.evaluations.BulkUploadError(cause: EvaluationAPIError, *, completed: BulkUploadReceipt, pending_results: tuple[EvaluationResult, ...], failed_index: int, failed_body: bytes) +class posthog.ai.evaluations.BulkUploadReceipt(chunks: tuple[UploadReceipt, ...]) +class posthog.ai.evaluations.CategoricalPassingRule +class posthog.ai.evaluations.CategoricalScorerConfig +class posthog.ai.evaluations.CategoricalScorerOption +class posthog.ai.evaluations.EvaluationAPIError(*, status: int | None = None, code: str | None = None, detail: Any = None, attr: str | None = None, errors: Any = None, response: dict[str, Any] | None = None, retry_after: float | None = None, persistence: Literal['rejected', 'unknown', 'not_sent'] = 'unknown', submission: dict[str, Any] | None = None) +class posthog.ai.evaluations.EvaluationItem(*, id: str | UUID | None = None, input: JSONValue | _Unset = _UNSET, output: JSONValue | _Unset = _UNSET, expected_output: JSONValue | _Unset = _UNSET, metadata: dict[str, JSONValue] | None | _Unset = _UNSET, case_key: str | None = None, trial: str | None = None, dataset_item_identifier: str | None = None, dataset_item_version_identifier: str | None = None, dataset_item_version_id: str | UUID | None = None, application_trace_id: str | None = None, payload: dict[str, JSONValue] | _Unset = _UNSET) +class posthog.ai.evaluations.EvaluationResult(*, item: EvaluationItem | str | UUID, scorer_version_id: str | UUID, value: ScoreValue | None = None, status: ResultStatus = 'ok', error_code: str | None = None, reasoning: str | None | _Unset = _UNSET, error_message: str | None | _Unset = _UNSET, metadata: dict[str, JSONValue] | None | _Unset = _UNSET, evaluator_trace_id: str | None = None, evaluated_at: datetime | str | None = None, payload: dict[str, JSONValue] | _Unset = _UNSET) +class posthog.ai.evaluations.Experiment(transport: SyncTransport, project_id: int, id: str, *, receipt: ExperimentReceipt | None = None, submission: bytes | None = None) +class posthog.ai.evaluations.ExperimentOptions +class posthog.ai.evaluations.ExperimentReceipt(id: str, status: Literal['uploading', 'completed', 'failed'], created: bool, started_at: datetime, created_at: datetime, finished_at: datetime | None, expected_item_count: int | None, expected_result_count: int | None, accepted_item_count: int, accepted_result_count: int) +class posthog.ai.evaluations.ItemReceipt(id: str, created: bool, accepted_at: datetime) +class posthog.ai.evaluations.NumericPassingRule +class posthog.ai.evaluations.NumericScorerConfig +class posthog.ai.evaluations.OfflineEvaluations(*, project_id: int, secret_key: str, host: str | None = None, timeout: float = 15, max_retries: int = 3) +class posthog.ai.evaluations.ResultOptions +class posthog.ai.evaluations.ResultReceipt(id: str, created: bool, accepted_at: datetime, item_id: str, scorer_version_id: str) +class posthog.ai.evaluations.Scorer(id: str, name: str, description: str, kind: ScorerKind, archived: bool, current_version: int, current_version_id: str, config: _ConfigT, created_at: datetime, updated_at: datetime, created_by: dict[str, Any] | None, team: int) +class posthog.ai.evaluations.ScorerPage(count: int, results: tuple[Scorer[ScorerConfig], ...], next: str | None, previous: str | None, next_offset: int | None) +class posthog.ai.evaluations.ScorerVersion(id: str, definition_id: str, version: int, kind: ScorerKind, config: _ConfigT, created_at: datetime, created_by: dict[str, Any] | None) +class posthog.ai.evaluations.ScorerVersionPage(count: int, results: tuple[ScorerVersion[ScorerConfig], ...], next_cursor: str | None) +class posthog.ai.evaluations.Scorers(transport: SyncTransport, project_id: int) +class posthog.ai.evaluations.UploadReceipt(items: tuple[ItemReceipt, ...], results: tuple[ResultReceipt, ...]) class posthog.ai.gemini.gemini.Client(api_key: Optional[str] = None, vertexai: Optional[bool] = None, credentials: Optional[Any] = None, project: Optional[str] = None, location: Optional[str] = None, debug_config: Optional[Any] = None, http_options: Optional[Any] = None, posthog_client: Optional[PostHogClient] = None, posthog_distinct_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: bool = False, posthog_groups: Optional[Dict[str, Any]] = None, **kwargs) class posthog.ai.gemini.gemini.Models(api_key: Optional[str] = None, vertexai: Optional[bool] = None, credentials: Optional[Any] = None, project: Optional[str] = None, location: Optional[str] = None, debug_config: Optional[Any] = None, http_options: Optional[Any] = None, posthog_client: Optional[PostHogClient] = None, posthog_distinct_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: bool = False, posthog_groups: Optional[Dict[str, Any]] = None, provider_client: Optional[Any] = None, **kwargs) class posthog.ai.gemini.gemini_async.AsyncClient(api_key: Optional[str] = None, vertexai: Optional[bool] = None, credentials: Optional[Any] = None, project: Optional[str] = None, location: Optional[str] = None, debug_config: Optional[Any] = None, http_options: Optional[Any] = None, posthog_client: Optional[PostHogClient] = None, posthog_distinct_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: bool = False, posthog_groups: Optional[Dict[str, Any]] = None, **kwargs) @@ -1317,6 +1460,38 @@ method posthog.ai.claude_agent_sdk.client.PostHogClaudeSDKClient.receive_respons method posthog.ai.claude_agent_sdk.client.PostHogClaudeSDKClient.set_model(model: Optional[str] = None) -> None method posthog.ai.claude_agent_sdk.client.PostHogClaudeSDKClient.set_permission_mode(mode: str) -> None method posthog.ai.claude_agent_sdk.processor.PostHogClaudeAgentProcessor.query(*, prompt: Any, options: Optional[ClaudeAgentOptions] = None, transport: Any = None, posthog_distinct_id: Optional[Union[str, Callable[[ResultMessage], Optional[str]]]] = None, posthog_trace_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: Optional[bool] = None, posthog_groups: Optional[Dict[str, Any]] = None) +method posthog.ai.evaluations.AsyncExperiment.complete() +method posthog.ai.evaluations.AsyncExperiment.fail() +method posthog.ai.evaluations.AsyncExperiment.upload_result(*, item: EvaluationItem | str | UUID, scorer_version_id: str | UUID, **options: Unpack[ResultOptions]) +method posthog.ai.evaluations.AsyncExperiment.upload_results(results: Sequence[EvaluationResult]) +method posthog.ai.evaluations.AsyncOfflineEvaluations.aclose() +method posthog.ai.evaluations.AsyncOfflineEvaluations.create_experiment(*, name: str, **options: Unpack[ExperimentOptions]) +method posthog.ai.evaluations.AsyncOfflineEvaluations.resume_experiment(id: str | UUID) +method posthog.ai.evaluations.AsyncScorers.create(*, name: str, kind: ScorerKind, config: ScorerConfig, description: str | None | _Omitted = _OMITTED) +method posthog.ai.evaluations.AsyncScorers.create_version(scorer_id: str | UUID, *, config: _ConfigT, base_version: int | None = None, name: str | _Omitted = _OMITTED, description: str | None | _Omitted = _OMITTED) +method posthog.ai.evaluations.AsyncScorers.get(scorer_id: str | UUID) +method posthog.ai.evaluations.AsyncScorers.get_version(scorer_id: str | UUID, version_id: str | UUID) +method posthog.ai.evaluations.AsyncScorers.list(*, limit: int = 100, offset: int = 0, search: str | None = None, kind: ScorerKind | None = None, archived: bool | None = None, order_by: ScorerOrder | None = None) +method posthog.ai.evaluations.AsyncScorers.list_versions(scorer_id: str | UUID, *, limit: int = 50, cursor: str | None = None) +method posthog.ai.evaluations.AsyncScorers.update(scorer_id: str | UUID, *, name: str | _Omitted = _OMITTED, description: str | None | _Omitted = _OMITTED, archived: bool | _Omitted = _OMITTED) +method posthog.ai.evaluations.EvaluationItem.from_dict(data: Mapping[str, Any]) +method posthog.ai.evaluations.EvaluationItem.to_dict() +method posthog.ai.evaluations.EvaluationResult.from_dict(data: Mapping[str, Any], *, item: EvaluationItem | None = None) +method posthog.ai.evaluations.EvaluationResult.to_dict() +method posthog.ai.evaluations.Experiment.complete() +method posthog.ai.evaluations.Experiment.fail() +method posthog.ai.evaluations.Experiment.upload_result(*, item: EvaluationItem | str | UUID, scorer_version_id: str | UUID, **options: Unpack[ResultOptions]) +method posthog.ai.evaluations.Experiment.upload_results(results: Sequence[EvaluationResult]) +method posthog.ai.evaluations.OfflineEvaluations.close() +method posthog.ai.evaluations.OfflineEvaluations.create_experiment(*, name: str, **options: Unpack[ExperimentOptions]) +method posthog.ai.evaluations.OfflineEvaluations.resume_experiment(id: str | UUID) +method posthog.ai.evaluations.Scorers.create(*, name: str, kind: ScorerKind, config: ScorerConfig, description: str | None | _Omitted = _OMITTED) +method posthog.ai.evaluations.Scorers.create_version(scorer_id: str | UUID, *, config: _ConfigT, base_version: int | None = None, name: str | _Omitted = _OMITTED, description: str | None | _Omitted = _OMITTED) +method posthog.ai.evaluations.Scorers.get(scorer_id: str | UUID) +method posthog.ai.evaluations.Scorers.get_version(scorer_id: str | UUID, version_id: str | UUID) +method posthog.ai.evaluations.Scorers.list(*, limit: int = 100, offset: int = 0, search: str | None = None, kind: ScorerKind | None = None, archived: bool | None = None, order_by: ScorerOrder | None = None) +method posthog.ai.evaluations.Scorers.list_versions(scorer_id: str | UUID, *, limit: int = 50, cursor: str | None = None) +method posthog.ai.evaluations.Scorers.update(scorer_id: str | UUID, *, name: str | _Omitted = _OMITTED, description: str | None | _Omitted = _OMITTED, archived: bool | _Omitted = _OMITTED) method posthog.ai.gemini.gemini.Models.embed_content(model: str, contents, posthog_distinct_id: Optional[str] = None, posthog_trace_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: Optional[bool] = None, posthog_groups: Optional[Dict[str, Any]] = None, **kwargs: Any) method posthog.ai.gemini.gemini.Models.generate_content(model: str, contents, posthog_distinct_id: Optional[str] = None, posthog_trace_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: Optional[bool] = None, posthog_groups: Optional[Dict[str, Any]] = None, **kwargs: Any) method posthog.ai.gemini.gemini.Models.generate_content_stream(model: str, contents, posthog_distinct_id: Optional[str] = None, posthog_trace_id: Optional[str] = None, posthog_properties: Optional[Dict[str, Any]] = None, posthog_privacy_mode: Optional[bool] = None, posthog_groups: Optional[Dict[str, Any]] = None, **kwargs: Any) @@ -1544,6 +1719,7 @@ module posthog.ai.claude_agent_sdk module posthog.ai.claude_agent_sdk.client module posthog.ai.claude_agent_sdk.formatting module posthog.ai.claude_agent_sdk.processor +module posthog.ai.evaluations module posthog.ai.gateway module posthog.ai.gemini module posthog.ai.gemini.gemini diff --git a/typings/requests/__init__.pyi b/typings/requests/__init__.pyi index 75a3fa48c..9e7892beb 100644 --- a/typings/requests/__init__.pyi +++ b/typings/requests/__init__.pyi @@ -1,6 +1,9 @@ from typing import Any -from . import adapters as adapters, exceptions as exceptions +from . import adapters as adapters, auth as auth, exceptions as exceptions + +class PreparedRequest: + headers: dict[str, str] class Response: status_code: int @@ -11,8 +14,20 @@ class Response: def close(self) -> None: ... class Session: + auth: auth.AuthBase | None + headers: dict[str, str] def mount(self, prefix: str, adapter: adapters.HTTPAdapter) -> None: ... def close(self) -> None: ... + def request( + self, + method: str, + url: str, + *, + data: str | bytes | None = ..., + params: dict[str, str | int] | None = ..., + timeout: float | None = ..., + allow_redirects: bool = ..., + ) -> Response: ... def post( self, url: str, diff --git a/typings/requests/auth.pyi b/typings/requests/auth.pyi new file mode 100644 index 000000000..5351e8390 --- /dev/null +++ b/typings/requests/auth.pyi @@ -0,0 +1,4 @@ +from . import PreparedRequest + +class AuthBase: + def __call__(self, request: PreparedRequest) -> PreparedRequest: ...