From a21e2a31827b305d7815cf4bfeab24855550b7a7 Mon Sep 17 00:00:00 2001 From: thomashebrard Date: Thu, 1 Oct 2026 13:50:57 +0200 Subject: [PATCH] Light run-list rows and result artifact selection list_runs and iterate_runs now yield RunHistoryItem rows carrying only the history fields the hosted API sends on GET /v1/runs (pipeline_run_id, status, created_at, finished_at, pipe_code, error). PipelineRun stays as the whole record behind RunDetail; the never-sent pipe_statuses map and its PipeStatus enum are removed. get_run_result, wait_for_result and start_and_wait take artifacts=, a sequence of the new RunArtifact enum, sent as one comma-separated ?artifacts= parameter. RunResults.main_stuff becomes optional and RunResults.carries(artifact) tells a not-requested artifact from a requested-but-unwritten one; the MissingMainStuffError check fires only when main_stuff was asked for. download_artifacts by run_id requests only the artifact its scope walks. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 15 +++ README.md | 4 +- docs/architecture.md | 12 +-- docs/artifact-download.md | 2 +- docs/run-results.md | 24 ++++- pipelex_sdk/artifact_models.py | 10 ++ pipelex_sdk/artifacts.py | 16 +-- pipelex_sdk/client.py | 81 ++++++++++++--- pipelex_sdk/product_models.py | 36 ++++--- pipelex_sdk/runs.py | 96 +++++++++++++----- tests/unit/test_artifacts.py | 26 ++++- tests/unit/test_client_lifecycle.py | 147 ++++++++++++++++++++++++++++ tests/unit/test_client_paging.py | 6 +- tests/unit/test_client_product.py | 17 +++- 14 files changed, 415 insertions(+), 77 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b74b789..ee76ae0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,20 @@ # Changelog +## [Unreleased] + +### Added + +- **Reading only some of a run's artifacts**: `get_run_result`, `wait_for_result` and `start_and_wait` take `artifacts=`, a sequence of the new `RunArtifact` enum, sent as one comma-separated `?artifacts=` parameter so the hosted API reads and re-signs only those artifacts; `None` still reads them all, and an empty selection raises `PipelineRequestError` before any request. `RunResults.carries(artifact)` tells an artifact the read did not ask for (absent) from one it asked for that was never written (`None`), and `download_artifacts` by `run_id` now asks only for the artifact its scope walks. See `docs/run-results.md`. + +### Changed + +- **Run-list rows are `RunHistoryItem` (Breaking)**: `list_runs` and `iterate_runs` yield `RunHistoryItem` rows carrying only `pipeline_run_id`, `status`, `created_at`, `finished_at`, `pipe_code` and `error`, which is all the hosted API now sends on `GET /v1/runs`; read `org_id`, `created_by_user_id`, `method_id`, `workflow_id` and `result_url` from `get_run_detail` instead. `PipelineRun` stays as the whole record and the base of `RunDetail`. +- **`RunResults.main_stuff` is optional (Breaking)**: it defaults to `None` so a selection without `MAIN_STUFF` parses; a read that asks for the main stuff (no selection, or one naming it) still raises `MissingMainStuffError` when a completed run delivers none. + +### Removed + +- **`PipeStatus` and `PipelineRun.pipe_statuses` (Breaking)**: the hosted API never sent a per-pipe status map on a run record, so the field and its enum are gone. + ## [v0.14.0] - 2026-09-27 ### Changed diff --git a/README.md b/README.md index 3afbcc7..ae385c2 100644 --- a/README.md +++ b/README.md @@ -165,7 +165,7 @@ except RunFailedError as exc: ... # the same run may succeed if started again ``` -Branch on `error_domain` (`input`, `config`, `runtime`), `type_uri` and `retryable`, never on the wording of `message`. The report is the runner's verbose one, so `message` and `provider_metadata` can hold a model provider's raw text: what a person should see of it is your application's decision. The same report is on `RunRead.error` when you read the run's status, on `RunResultFailed.error` from `get_run_result`, and on `PipelineRun.error` in the run lists. +Branch on `error_domain` (`input`, `config`, `runtime`), `type_uri` and `retryable`, never on the wording of `message`. The report is the runner's verbose one, so `message` and `provider_metadata` can hold a model provider's raw text: what a person should see of it is your application's decision. The same report is on `RunRead.error` when you read the run's status, on `RunResultFailed.error` from `get_run_result`, and on `RunHistoryItem.error` in the run lists. ### API errors: branch on `type_uri` and `error_domain`, not the HTTP status @@ -196,7 +196,7 @@ There is no barrel import — package `__init__.py` files stay empty. Import eac - **Client & construction** — `from pipelex_sdk.client import PipelexAPIClient, DEFAULT_API_BASE_URL, MthdsFile` - **Run lifecycle types** — `from pipelex_sdk.runs import RunStatus, RunPublic, RunRead, RunResults, RunResultState, WaitForResultOptions, PollInfo` - **Error reports** — `from pipelex_sdk.error_models import RunErrorReport, UserAction, ProviderErrorMetadata, MigrationErrorBlock, FieldError` -- **Product wire models** — `from pipelex_sdk.product_models import UserProfile, MethodData, MethodWriteInput, Membership, MembershipsResponse, SubscriptionResponse, PlanView, InvoiceView, OnboardingSubmission, UploadInput, UploadedFile, PipelineRun, ...`, with the catalog-source readers beside them: `method_source_to_contents` turns a fetched `MethodData.mthds` into the `mthds_contents` a run or a validate takes, and `MethodFile` / `parse_method_files` / `serialize_method_files` are the codec for a method's custom PipeFunc `python`. +- **Product wire models** — `from pipelex_sdk.product_models import UserProfile, MethodData, MethodWriteInput, Membership, MembershipsResponse, SubscriptionResponse, PlanView, InvoiceView, OnboardingSubmission, UploadInput, UploadedFile, RunHistoryItem, RunDetail, ...`, with the catalog-source readers beside them: `method_source_to_contents` turns a fetched `MethodData.mthds` into the `mthds_contents` a run or a validate takes, and `MethodFile` / `parse_method_files` / `serialize_method_files` are the codec for a method's custom PipeFunc `python`. - **Validation verdict types** — `from pipelex_sdk.validation_models import PipelexValidationResult, PipelexValidationReport, PipelexInvalidReport, ValidationErrorItem, SuggestedFix, VALIDATION_VIEW_INPUT_FORM, ...` - **Codegen tree** — `from pipelex_sdk.codegen_writer import write_codegen_tree, CodegenTreeWriteReport` to write one, `from pipelex_sdk.codegen_check import run_codegen_check, CodegenCheckReport, CodegenDrift, DriftCategory` to verify one, with the format primitives in `pipelex_sdk.codegen_lock` (`CodegenLock`, `parse_lock`, `load_lock`, `validate_artifact_path`, ...) and `pipelex_sdk.codegen_stamp` (`STAMPABLE_SUFFIXES`, `is_stampable_artifact_path`, `compute_content_hash`, `parse_stamped`, ...) - **Typed errors** — `from pipelex_sdk.errors import ApiResponseError, ApiUnreachableError, PipelineExecuteTimeoutError, PagingNotTerminatingError, RunFailedError, RunTimeoutError, RunLifecycleUnavailableError, RunStillRunningError, CodegenError, CodegenLockError, ...` diff --git a/docs/architecture.md b/docs/architecture.md index 2e43f90..86ba180 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -121,8 +121,8 @@ The durable run lifecycle (`pipelex_sdk/runs.py` + the client's lifecycle method ### Polling surface - **`get_run_status(run_id)`** — `GET /v1/runs/{id}/status` → `RunRead`. Lifts the `Retry-After` header onto `retry_after_seconds`. -- **`get_run_result(run_id)`** — `GET /v1/runs/{id}/results`, mapping the platform's poll semantics to the `RunResultState` union: `202`/`503` → `running` (in-flight / degraded — never fail a poller), `200` → `completed`, `409` → `failed`. The `409`'s problem document carries `detail` (`Run finished with status : `), which becomes `message`, and two extension members: `run_status`, which becomes `status`, and `error`, the run's stored report, which becomes `error` typed as `RunErrorReport`. The status is read from that member and never parsed out of the sentence; a `409` without a `run_status` this SDK knows reads as `FAILED`. The report is read leniently, because a report written by another runner version must not mask the failure it explains: a known field whose value does not fit its type reads as `None` and the rest of the report stands, and only an `error` that is not an object at all reads as `None` — `message` still carries the reason either way. The status read and the run lists read it the same way (`LenientRunErrorReport`). -- **`wait_for_result(run_id, options)`** — polls `get_run_result` to a terminal state, honoring `Retry-After` and the deadline. Resolves on `COMPLETED`; raises `RunFailedError` on any other terminal status, carrying the failed arm's `status`, `message` and `error`, and `RunTimeoutError` if the budget elapses (the run keeps executing server-side — resume later by id). +- **`get_run_result(run_id, *, artifacts=None)`** — `GET /v1/runs/{id}/results`, with `artifacts` (a sequence of `RunArtifact`) sent as one comma-separated `?artifacts=` parameter that narrows the `200` body to the named artifacts; an unselected artifact is absent from the result, which `RunResults.carries(artifact)` reports, and the `MissingMainStuffError` check runs only when `main_stuff` was asked for ([run-results.md](./run-results.md#reading-only-some-artifacts)). It maps the platform's poll semantics to the `RunResultState` union: `202`/`503` → `running` (in-flight / degraded — never fail a poller), `200` → `completed`, `409` → `failed`. The `409`'s problem document carries `detail` (`Run finished with status : `), which becomes `message`, and two extension members: `run_status`, which becomes `status`, and `error`, the run's stored report, which becomes `error` typed as `RunErrorReport`. The status is read from that member and never parsed out of the sentence; a `409` without a `run_status` this SDK knows reads as `FAILED`. The report is read leniently, because a report written by another runner version must not mask the failure it explains: a known field whose value does not fit its type reads as `None` and the rest of the report stands, and only an `error` that is not an object at all reads as `None` — `message` still carries the reason either way. The status read and the run lists read it the same way (`LenientRunErrorReport`). +- **`wait_for_result(run_id, options, *, artifacts=None)`** — polls `get_run_result` with the same selection to a terminal state, honoring `Retry-After` and the deadline. Resolves on `COMPLETED`; raises `RunFailedError` on any other terminal status, carrying the failed arm's `status`, `message` and `error`, and `RunTimeoutError` if the budget elapses (the run keeps executing server-side — resume later by id). These poll GETs go through `_send_or_unreachable`, so a transport failure surfaces as `ApiUnreachableError` (consistent with the product layer), while a missing-route `404` surfaces as `RunLifecycleUnavailableError` and any other non-2xx — the platform's run-not-found `404` included — as `ApiResponseError`. @@ -225,7 +225,7 @@ Beyond those two the verdicts match, including the drift sentences. One differen The hosted catalog/account routes the webapp drives (`pipelex_sdk/product_models.py` + the client's product methods). Every route rides the same `{base}/v1/*` surface, `Authorization: Bearer`, org-from-JWT contract as the protocol routes, and goes through `_request_product`, which maps a non-2xx `problem+json` to a typed `ApiResponseError` — **consumers branch on `.type_uri` (and on `.error_domain` where a runner-rendered problem carries it), never the HTTP status**, and read `.code` for the platform's native code (see Error regimes above). -The wire models are snake_case Pydantic v2. Response models are extension-open (`extra="allow"`) so a newly-added server field is preserved, not rejected; input models name exactly what each route accepts. `PipelineRun.status` reuses the run-lifecycle `RunStatus`; `OrgRole`, `PipeStatus`, and the onboarding fields are `StrEnum`s. +The wire models are snake_case Pydantic v2. Response models are extension-open (`extra="allow"`) so a newly-added server field is preserved, not rejected; input models name exactly what each route accepts. `RunHistoryItem.status` and `PipelineRun.status` reuse the run-lifecycle `RunStatus`; `OrgRole` and the onboarding fields are `StrEnum`s. - **User profile** — `get_me()` → `UserProfile` (`GET /v1/me`). - **Methods catalog** — `list_methods()` / `iterate_methods()` / `get_method(id)` / `create_method(MethodWriteInput)` / `update_method(id, MethodWriteInput)` (a rename is a changed `name`) / `delete_method(id)`. The id is path-encoded; an absent `input_data` is dropped from the write body. @@ -250,11 +250,11 @@ The wire models are snake_case Pydantic v2. Response models are extension-open ( - **Storage** — `resolve_storage_url(uri)` → presigned URL for one reference; `resolve_storage_urls_bulk(uris)` → one verdict per reference through `POST /v1/resolve-storage-url/bulk`, the route the artifact stack is built on; `upload(UploadInput)` → the stored file handle. The higher-level `upload_file` / `prepare_inputs` preparation surface built on top of `upload` is now available — see [input-preparation.md](./input-preparation.md) — and its download twin is the artifact stack below. - **Run records** — `list_runs(method_id, …)` / `iterate_runs(method_id, …)` / `get_run_detail(run_id)` (the catalog-style reads, distinct from the lifecycle status/result routes); `update_run(run_id, UpdateRunInput)` (admin/manual status patch, empty 2xx body). - **This list is paged too.** `list_runs(method_id, created_from=…, created_to=…, limit=…, cursor=…)` returns a `RunPage` with the same opaque-cursor contract as `MethodPage`. `created_from` / `created_to` are **instants** — ISO-8601 with a UTC offset — and inclusive; they are index key conditions rather than filters, so a bare date or a naive timestamp is a platform `400` surfaced as `ApiResponseError`. Worth knowing and not obvious from the route: every `/v1/runs*` product route sits behind the platform's surface-access gate, which for API-key auth demands the `ff_api_keys` feature flag and fails closed with a `403` — so a `403` here means "flag", not "wrong key". + **This list is paged too.** `list_runs(method_id, created_from=…, created_to=…, limit=…, cursor=…)` returns a `RunPage` of `RunHistoryItem` rows with the same opaque-cursor contract as `MethodPage`. A row is what a history shows and nothing more: `pipeline_run_id`, `status`, `created_at`, `finished_at`, `pipe_code` and `error`. The organization, the creator, the method (the caller named it), the workflow id and the result prefix are not on the list; `get_run_detail` serves the whole record when a run is opened. `created_from` / `created_to` are **instants** — ISO-8601 with a UTC offset — and inclusive; they are index key conditions rather than filters, so a bare date or a naive timestamp is a platform `400` surfaced as `ApiResponseError`. Worth knowing and not obvious from the route: every `/v1/runs*` product route sits behind the platform's surface-access gate, which for API-key auth demands the `ff_api_keys` feature flag and fails closed with a `403` — so a `403` here means "flag", not "wrong key". **The two iterators stop on different signals, and the difference is in the server.** `iterate_methods` continues through an empty page with a live cursor; `iterate_runs` treats an **empty page as the end**, because the run date bounds are index key conditions and so a run page is never empty-with-a-cursor. Both share the same runaway page ceiling, and both raise `PagingNotTerminatingError` at it: the empty-page stop only catches a server minting fresh cursors while returning *nothing*, so a cursor that cycles across two or more values (`c1 → c2 → c1`) over non-empty pages trips neither that check nor the adjacent-cursor one. The ceiling is the cheap guard against that whole family — tracking every cursor seen would cost unbounded memory for the same protection. - **`PipelineRun` fields the platform genuinely serves as null are typed nullable.** `method_id` is `None` for an ad-hoc run from an inline bundle, which belongs to no stored method; `pipe_code` is `None` for a run that let the bundle's `main_pipe` decide. `org_id`, `created_by_user_id`, and `error: RunErrorReport | None` — the run's stored report, typed whole exactly as on the status read — join them. `RunDetail`, returned only by `get_run_detail`, adds `mthds_contents` and `inputs`: what the run actually executed, and the only record of it, since a method edited since the run no longer describes what happened. Both are left out of the list and the polled status on purpose — their cost scales with page size and poll rate respectively. + **Fields the platform genuinely serves as null are typed nullable.** On `RunHistoryItem`, `pipe_code` is `None` for a run that let the bundle's `main_pipe` decide, and `error` is `None` on every run that did not fail. On `PipelineRun`, the whole record and the base of `RunDetail`, `method_id` is `None` for an ad-hoc run from an inline bundle, which belongs to no stored method; `pipe_code` is `None` for a run that let the bundle's `main_pipe` decide. `org_id`, `created_by_user_id`, and `error: RunErrorReport | None` — the run's stored report, typed whole exactly as on the status read — join them. `RunDetail`, returned only by `get_run_detail`, adds `mthds_contents` and `inputs`: what the run actually executed, and the only record of it, since a method edited since the run no longer describes what happened. Both are left out of the list and the polled status on purpose — their cost scales with page size and poll rate respectively. ## Artifact stack (`pipelex_sdk/artifacts.py` + `pipelex_sdk/artifact_models.py`) @@ -293,7 +293,7 @@ Each stays deferred rather than silently missing. Everything else — the protoc **The run-results surface** — `RunResults` tracks `@pipelex/sdk` 0.20.0 field for field, and is at parity with it. Every field is declared, with the two-path semantics [`run-results.md`](./run-results.md) states: `graph_spec` (lifted on the blocking path instead of written `None`), `graph_assembly_error`, the three I/O artifacts typed from `mthds.protocol` with `pipe_io_artifacts_error` beside them, the usage pair, `working_memory` (the hosted artifact relayed as its own key, the blocking one lifted off `pipe_output.working_memory`), and `pipe_output`; and the usage pair's fold, `summarize_usage`, with the summary types it returns. What stays different is idiomatic and deliberate: a key the hosted body did not carry is absent from `model_fields_set` where the JS reads `undefined`, and the page says how to read that. -**Models** — field-for-field across the run-lifecycle types and the product wire models. Deliberate idiomatic ports (not gaps): milliseconds → seconds (`interval_seconds` / `timeout_seconds` / `elapsed_seconds`); the JS `AbortSignal` → Python `asyncio` cancellation (no `signal` field, and no `aborted` flag on the download verdict — a cancelled `download_artifacts` raises `CancelledError` with its partial files already unlinked); the JS artifact request object → `DownloadArtifactsOptions` beside an explicit `dir_path` (`dir` being a Python builtin) and the raw bulk call named `resolve_storage_urls_bulk` after its route rather than `resolveStorageUrls`, so it cannot be misread as the single-reference `resolve_storage_url`; the returned `Response` of `fetchArtifact` → an async context manager, because an httpx stream is only live inside its own block; JS inline string-unions promoted to `StrEnum`s (`OrgRole`, `PipeStatus`, the onboarding fields) with identical wire values; response models are `extra="allow"` for forward-compat. The Pipelex validation narrowing is **owned here** (`pipelex_sdk.validation_models`), narrowing `mthds`'s neutral verdict bases (the resolved follow-up #9); the brand-neutral `Dict*` wire concretes (`DictRunResultExecute`, and the `DictPipeOutputAbstract` / `DictWorkingMemoryAbstract` pair that types `RunResults.pipe_output` and `RunResults.working_memory`) are reused from `mthds` by inheritance or by import — they are a shared wire contract the `pipelex` runtime also builds on — rather than duplicated as `pipelex-sdk-js` does. One addition runs ahead of the JS SDK: `write_codegen_tree` has no `@pipelex/sdk` counterpart, because the JS writer lives inside `pipelex-starter-js`'s own harness, mixed in with project policy; the byte-fidelity contract is small and load-bearing enough to belong to the SDK, where every consumer shares one correct implementation. The drift check goes the other way and takes a different shape on purpose: `@pipelex/sdk`'s `runCodegenCheck` is **pure** — the caller walks its own tree and hands in the text — because that module must stay free of Node builtins for a browser bundle, and its doc consequently loads the caller with obligations (walk the whole tree, do not reformat, decode strictly) whose every breach yields a wrong verdict rather than an error. Python has no such constraint, this package already does filesystem work, and `pipelex`'s own surface takes a root — so `run_codegen_check(root=…)` takes one too. Every caller obligation becomes the library's, the verdict is provable against `pipelex` by calling both with the same directory, and it composes with `write_codegen_tree(report, output_dir=…)` as the same path in and out. Two divergences worth naming: the page envelopes keep the wire's snake_case `next_cursor`, where the JS mirror renamed it `nextCursor` for its own consumers; and the method-files catalog converter (`parse_method_files` / `serialize_method_files`) lives in this package, where the JS pair lives in `mthds-js` because `pipelex-mcp` consumes the same format and wanted one owner. There is no second Python consumer, and the catalog serialization is a Pipelex product concern rather than an MTHDS protocol one, so this SDK is a proper home for it. `mthds` has since grown the canonical pair as `mthds.protocol.method_files`, shipped since `mthds` 0.15.0 and present in the `mthds==0.16.0` this package pins exactly, and the local pair is kept beside it on purpose rather than left over. Adoption is a step of its own rather than a rename: that parser raises `PipelineRequestError`, which does not subclass `ValueError`, and pydantic converts only `ValueError` out of a validator — so re-pointing `MethodData`'s validator at it would stop every caller's `except ValidationError` from catching a malformed stored source. The exception's base class is `mthds`'s to settle, and it is settled first; until then the two implementations are kept in step, and they disagree on blankness (this one is Python's `str.strip`, the canonical one is ECMAScript's), on the serialized bytes (default separators and ASCII escapes here, `JSON.stringify`'s there), and on the JSON constants and integer-literal cap the canonical one closes with `parse_constant=` and `parse_int=float`. +**Models** — field-for-field across the run-lifecycle types and the product wire models. Deliberate idiomatic ports (not gaps): milliseconds → seconds (`interval_seconds` / `timeout_seconds` / `elapsed_seconds`); the JS `AbortSignal` → Python `asyncio` cancellation (no `signal` field, and no `aborted` flag on the download verdict — a cancelled `download_artifacts` raises `CancelledError` with its partial files already unlinked); the JS artifact request object → `DownloadArtifactsOptions` beside an explicit `dir_path` (`dir` being a Python builtin) and the raw bulk call named `resolve_storage_urls_bulk` after its route rather than `resolveStorageUrls`, so it cannot be misread as the single-reference `resolve_storage_url`; the returned `Response` of `fetchArtifact` → an async context manager, because an httpx stream is only live inside its own block; JS inline string-unions promoted to `StrEnum`s (`OrgRole`, the onboarding fields) with identical wire values; response models are `extra="allow"` for forward-compat. The Pipelex validation narrowing is **owned here** (`pipelex_sdk.validation_models`), narrowing `mthds`'s neutral verdict bases (the resolved follow-up #9); the brand-neutral `Dict*` wire concretes (`DictRunResultExecute`, and the `DictPipeOutputAbstract` / `DictWorkingMemoryAbstract` pair that types `RunResults.pipe_output` and `RunResults.working_memory`) are reused from `mthds` by inheritance or by import — they are a shared wire contract the `pipelex` runtime also builds on — rather than duplicated as `pipelex-sdk-js` does. One addition runs ahead of the JS SDK: `write_codegen_tree` has no `@pipelex/sdk` counterpart, because the JS writer lives inside `pipelex-starter-js`'s own harness, mixed in with project policy; the byte-fidelity contract is small and load-bearing enough to belong to the SDK, where every consumer shares one correct implementation. The drift check goes the other way and takes a different shape on purpose: `@pipelex/sdk`'s `runCodegenCheck` is **pure** — the caller walks its own tree and hands in the text — because that module must stay free of Node builtins for a browser bundle, and its doc consequently loads the caller with obligations (walk the whole tree, do not reformat, decode strictly) whose every breach yields a wrong verdict rather than an error. Python has no such constraint, this package already does filesystem work, and `pipelex`'s own surface takes a root — so `run_codegen_check(root=…)` takes one too. Every caller obligation becomes the library's, the verdict is provable against `pipelex` by calling both with the same directory, and it composes with `write_codegen_tree(report, output_dir=…)` as the same path in and out. Two divergences worth naming: the page envelopes keep the wire's snake_case `next_cursor`, where the JS mirror renamed it `nextCursor` for its own consumers; and the method-files catalog converter (`parse_method_files` / `serialize_method_files`) lives in this package, where the JS pair lives in `mthds-js` because `pipelex-mcp` consumes the same format and wanted one owner. There is no second Python consumer, and the catalog serialization is a Pipelex product concern rather than an MTHDS protocol one, so this SDK is a proper home for it. `mthds` has since grown the canonical pair as `mthds.protocol.method_files`, shipped since `mthds` 0.15.0 and present in the `mthds==0.16.0` this package pins exactly, and the local pair is kept beside it on purpose rather than left over. Adoption is a step of its own rather than a rename: that parser raises `PipelineRequestError`, which does not subclass `ValueError`, and pydantic converts only `ValueError` out of a validator — so re-pointing `MethodData`'s validator at it would stop every caller's `except ValidationError` from catching a malformed stored source. The exception's base class is `mthds`'s to settle, and it is settled first; until then the two implementations are kept in step, and they disagree on blankness (this one is Python's `str.strip`, the canonical one is ECMAScript's), on the serialized bytes (default separators and ASCII escapes here, `JSON.stringify`'s there), and on the JSON constants and integer-literal cap the canonical one closes with `parse_constant=` and `parse_int=float`. **Errors** — `ApiUnreachableError`, `PipelineExecuteTimeoutError`, `RunFailedError`, `RunTimeoutError`, `RunLifecycleUnavailableError`, `PagingNotTerminatingError` are owned here; `ApiResponseError` is owned here as a narrowing of `mthds`'s own, with the same members as `@pipelex/sdk`'s apart from the inference members (`model`, `provider`, `provider_metadata`, `migration`) the JS error also lifts, which a Python caller reads off `problem`; `RunStillRunningError` is re-exported from `mthds`. `ClientAuthenticationError` is **not** ported: it is a dormant export in the JS barrel (defined and exported but never raised by the client), and in Python it already lives in `mthds.runners.api.exceptions` — importable directly if ever needed, with no barrel here to re-export it through. diff --git a/docs/artifact-download.md b/docs/artifact-download.md index 89b6514..947e9ab 100644 --- a/docs/artifact-download.md +++ b/docs/artifact-download.md @@ -97,7 +97,7 @@ if not verdict.all_saved: `client.download_artifacts(dir_path=…, run_id=…)` is the same thing as a method. -**Where it reads from.** Exactly one of `run_id` and `results`. A `run_id` re-reads the results through `get_run_result`, so a completed run is downloadable days later from its id alone; a `RunResults` already in hand is read as it is, with no request. `scope` picks the artifact walked for references: `main_stuff` (the default) is the run's output, and `working_memory` is the opt-in that also brings down the echoed inputs and every intermediate stuff — it is read off `RunResults.working_memory`, the declared field both paths deliver ([`run-results.md`](./run-results.md#working_memory--every-named-stuff-of-the-run)). +**Where it reads from.** Exactly one of `run_id` and `results`. A `run_id` re-reads the results through `get_run_result`, asking only for the artifact the scope walks, so a completed run is downloadable days later from its id alone; a `RunResults` already in hand is read as it is, with no request. `scope` picks the artifact walked for references: `main_stuff` (the default) is the run's output, and `working_memory` is the opt-in that also brings down the echoed inputs and every intermediate stuff — it is read off `RunResults.working_memory`, the declared field both paths deliver ([`run-results.md`](./run-results.md#working_memory--every-named-stuff-of-the-run)). **How it downloads.** The whole set is resolved through the bulk route ahead of the work, then an `asyncio.Semaphore` bounds how many references are in flight at once (`concurrency`, default 4) over the whole per-reference pipeline: fetch, create the file exclusively, stream the body in. Resolution is just-in-time where it matters: a link that has expired by the time its task reaches it — a large set downloaded a few at a time can outlive the fifteen-minute link — is resolved again for that reference alone, so no fetch ever runs on a stale signature. diff --git a/docs/run-results.md b/docs/run-results.md index e1584b3..c1addee 100644 --- a/docs/run-results.md +++ b/docs/run-results.md @@ -18,7 +18,7 @@ print(results.main_stuff, summarize_usage(results).total_cost_usd) | field | type | hosted (durable) path | bare-runner (blocking) path | |---|---|---|---| | `pipeline_run_id` | `str` | the run store's id | the runner's own id for the call | -| `main_stuff` | `Any` | the `main_stuff.json` artifact | resolved out of the returned working memory | +| `main_stuff` | `Any` | the `main_stuff.json` artifact (absent when a selection leaves it out) | resolved out of the returned working memory | | `graph_spec` | `Any` | the `graphspec.json` artifact | lifted off `pipe_output` | | `graph_assembly_error` | `str \| None` | absent until the platform relays it | lifted off `pipe_output` | | `pipe_io_contracts` | `PipeIOContracts \| None` | the `pipe_io_contracts.json` artifact | lifted off `pipe_output` | @@ -32,7 +32,7 @@ print(results.main_stuff, summarize_usage(results).total_cost_usd) ## `None` versus absent — the reading this page relies on -Every field but the first two is optional, and two readings of an optional field are distinct on purpose. The JS SDK reads a key the hosted body did not carry as `undefined` and a key relayed as `null` as `null`. Python has one `None`, so the SDK keeps the distinction where pydantic keeps it: in `model_fields_set`. A key the hosted body did not carry is not in the set and reads `None`; a key relayed as `null` is in the set and reads `None` too. +Every field but `pipeline_run_id` is optional, and two readings of an optional field are distinct on purpose. The JS SDK reads a key the hosted body did not carry as `undefined` and a key relayed as `null` as `null`. Python has one `None`, so the SDK keeps the distinction where pydantic keeps it: in `model_fields_set`. A key the hosted body did not carry is not in the set and reads `None`; a key relayed as `null` is in the set and reads `None` too. ```python results = await client.wait_for_result(run_id) @@ -45,6 +45,24 @@ elif results.graph_assembly_error is None: That distinction matters on the hosted path only. On the blocking path the SDK lifts every field off the runner's output and passes each one explicitly, so every field is set there whether or not the runner carried the key — the blocking path always answers, exactly as the JS twin writes `null` for each. Most consumers never need the set: a check of `is None` is the right branch for "is there a value", and `model_fields_set` is for the one question it answers, whether the wire said anything at all. +## Reading only some artifacts + +A hosted results read fetches every artifact from the run store and re-signs every link inside them, which is wasted work for a caller that wants one of them, such as a run history showing each run's output. `get_run_result`, `wait_for_result` and `start_and_wait` take `artifacts=`, a sequence of `RunArtifact` (`pipelex_sdk/runs.py`): `GRAPH_SPEC`, `PIPE_IO_CONTRACTS`, `INPUT_FORM`, `OUTPUT_FORM`, `MAIN_STUFF`, `WORKING_MEMORY` and `TOKENS_USAGES`, each named after the field it fills. The client sends the selection as one comma-separated `?artifacts=` parameter, deduplicated and in that declaration order, and the platform then reads, re-signs and returns only those. `TOKENS_USAGES` names the usage envelope, so it fills `usage_assembly_error` too. `None`, the default, reads everything, exactly as a read without the parameter always has; an empty selection names nothing and is refused before any request with `PipelineRequestError` (the platform would answer it with a `400`, as it does an unknown name). + +```python +from pipelex_sdk.runs import RunArtifact, RunResultCompleted + +state = await client.get_run_result(run_id, artifacts=[RunArtifact.MAIN_STUFF]) +if isinstance(state, RunResultCompleted): + print(state.result.main_stuff) # graph_spec, the forms, the memory and the usage were never read +``` + +**Not requested against requested but not written.** The platform leaves an unselected artifact out of the body and relays a selected one that was never written as `null`, which is the distinction the next section describes, read per artifact: `results.carries(RunArtifact.GRAPH_SPEC)` is `False` when the read did not ask for the graph, and `True` with `results.graph_spec is None` when it asked and the run store has none. `carries` checks every field the artifact fills, so `TOKENS_USAGES` is carried only when both `tokens_usages` and `usage_assembly_error` came back. + +**The main-stuff check follows the selection.** A completed run read with no selection, or with one naming `MAIN_STUFF`, must deliver a main stuff, and `MissingMainStuffError` is raised when it does not. A selection that leaves `MAIN_STUFF` out asked for none, so `results.main_stuff` reads `None`, `carries(RunArtifact.MAIN_STUFF)` is `False`, and nothing is raised. + +**The blocking path ignores the selection.** `start_and_wait` against a bare runner gets every artifact in the one execute response, so narrowing would save nothing; it returns the full result, every field set, whatever was asked for. `download_artifacts` by `run_id` asks for the one artifact its scope walks, `main_stuff` or `working_memory`, and nothing else. + ## `pipeline_run_id` — the durable handle The run id is what makes a run readable after the process that started it has gone. `start` returns it in its acknowledgement before the run finishes, and every lifecycle read takes it: `get_run_status(run_id)` for the status row, `get_run_result(run_id)` for a single result lookup, `wait_for_result(run_id)` to resume polling a run an earlier session started. It is also what a `RunTimeoutError` leaves you with — the run keeps executing server-side, so the timeout is a reason to re-poll by id, not a reason to run the method again. @@ -61,7 +79,7 @@ Against a bare runner the id identifies the call the runner just answered, but t ## `main_stuff` — the output -`main_stuff` is the resolved content of the run's main output and is always present for a completed run. On the hosted path it is the `main_stuff.json` artifact; on the blocking path the SDK resolves it out of the returned working memory through the response's `main_stuff_name`. Both deliver the same content shape, so there is no shape-guessing and no path-dependent branch to write. A completed run that cannot deliver one raises `MissingMainStuffError` rather than handing back a half-filled result. +`main_stuff` is the resolved content of the run's main output and is always present for a completed run read in full, or read with a selection that names it (a selection that leaves it out reads it as `None`, [above](#reading-only-some-artifacts)). On the hosted path it is the `main_stuff.json` artifact; on the blocking path the SDK resolves it out of the returned working memory through the response's `main_stuff_name`. Both deliver the same content shape, so there is no shape-guessing and no path-dependent branch to write. A completed run that cannot deliver one raises `MissingMainStuffError` rather than handing back a half-filled result. It is typed `Any` because the content is polymorphic: a structured output arrives as a dict of the concept's fields, and a multiple output as the envelope `{"items": [...]}` that the runtime's `ListContent` serialises to. Every content type serialises to an object, natives included — a text output is `{"text": "…"}` and a number `{"number": 0}` — so a guard written for a bare `""` or `0` never fires, and an empty multiple output is `{"items": []}` rather than `[]`. Narrow it where you read it, ideally through the types generated for the method rather than a hand-written cast. diff --git a/pipelex_sdk/artifact_models.py b/pipelex_sdk/artifact_models.py index d68e99b..49ae12c 100644 --- a/pipelex_sdk/artifact_models.py +++ b/pipelex_sdk/artifact_models.py @@ -15,6 +15,7 @@ from pydantic import BaseModel, ConfigDict, Field from pipelex_sdk._pydantic_utils import empty_list_factory_of +from pipelex_sdk.runs import RunArtifact # ── Constants ──────────────────────────────────────────────────────── @@ -59,6 +60,15 @@ def results_field(self) -> str: case ArtifactScope.WORKING_MEMORY: return "working_memory" + @property + def run_artifact(self) -> RunArtifact: + """The one result artifact this scope needs, so a download by run id reads nothing else.""" + match self: + case ArtifactScope.MAIN_STUFF: + return RunArtifact.MAIN_STUFF + case ArtifactScope.WORKING_MEMORY: + return RunArtifact.WORKING_MEMORY + class ArtifactItemError(BaseModel): """Why one reference failed — a value, never a raised error. diff --git a/pipelex_sdk/artifacts.py b/pipelex_sdk/artifacts.py index 835bc4a..22d1588 100644 --- a/pipelex_sdk/artifacts.py +++ b/pipelex_sdk/artifacts.py @@ -65,7 +65,7 @@ RunStillRunningError, ScopeUnavailableError, ) -from pipelex_sdk.runs import RunResultCompleted, RunResultFailed, RunResultRunning, RunResults +from pipelex_sdk.runs import RunArtifact, RunResultCompleted, RunResultFailed, RunResultRunning, RunResults if TYPE_CHECKING: from collections.abc import AsyncGenerator, AsyncIterator, Sequence @@ -154,7 +154,7 @@ async def resolve_storage_urls_bulk(self, uris: list[str]) -> BulkResolvedStorag class ArtifactCapableClient(BulkResolveClient, Protocol): """What `download_artifacts` needs on top: the single-shot result lookup, for the `run_id` arm.""" - async def get_run_result(self, run_id: str) -> RunResultState: ... + async def get_run_result(self, run_id: str, *, artifacts: Sequence[RunArtifact] | None = None) -> RunResultState: ... # ── locate_artifacts / collect_artifacts ───────────────────────────── @@ -785,7 +785,7 @@ async def download_artifacts( if bool(run_id) == (results is not None): msg = "download_artifacts takes exactly one of `run_id` (the results are re-read) or `results` (a RunResults in hand)." raise ArtifactOperationError(msg) - read_results = results if results is not None else await _read_completed_results(client, cast("str", run_id)) + read_results = results if results is not None else await _read_completed_results(client, cast("str", run_id), scope) walked = _scope_value(read_results, scope) # The walk's own record names the files; `locations` is what the verdict reports. @@ -843,9 +843,13 @@ async def download_artifacts( return verdict -async def _read_completed_results(client: ArtifactCapableClient, run_id: str) -> RunResults: - """Read a run's results by id, turning a run that has not completed into its typed error.""" - state = await client.get_run_result(run_id) +async def _read_completed_results(client: ArtifactCapableClient, run_id: str, scope: ArtifactScope) -> RunResults: + """Read a run's results by id, turning a run that has not completed into its typed error. + + Asks for the scope's one artifact only, so a download never pays for the graph, the forms or + the usage it does not walk. + """ + state = await client.get_run_result(run_id, artifacts=[scope.run_artifact]) if isinstance(state, RunResultRunning): retry = state.retry_after_seconds hint = f" — retry in {retry}s." if retry is not None else "." diff --git a/pipelex_sdk/client.py b/pipelex_sdk/client.py index b72262a..ee192bb 100644 --- a/pipelex_sdk/client.py +++ b/pipelex_sdk/client.py @@ -79,10 +79,10 @@ MethodSummary, PipelexApiKeyCreated, PipelexApiKeyList, - PipelineRun, PlanView, ResolvedStorageUrl, RunDetail, + RunHistoryItem, RunPage, SubscriptionResponse, UploadedFile, @@ -91,6 +91,7 @@ from pipelex_sdk.runs import ( PipelexRunResultStart, PollInfo, + RunArtifact, RunRead, RunResultCompleted, RunResultFailed, @@ -105,7 +106,7 @@ from pipelex_sdk.validation_models import PipelexValidationResultAdapter, ValidationErrorItem if TYPE_CHECKING: - from collections.abc import AsyncIterator + from collections.abc import AsyncIterator, Sequence from contextlib import AbstractAsyncContextManager from pathlib import Path @@ -795,7 +796,7 @@ async def get_run_status(self, run_id: str) -> RunRead: run = run.model_copy(update={"retry_after_seconds": retry_after}) return run - async def get_run_result(self, run_id: str) -> RunResultState: + async def get_run_result(self, run_id: str, *, artifacts: Sequence[RunArtifact] | None = None) -> RunResultState: """Single-shot result lookup — `GET /v1/runs/{run_id}/results`. Maps the platform's poll semantics to a discriminated union: @@ -805,12 +806,28 @@ async def get_run_result(self, run_id: str) -> RunResultState: - HTTP 409 → `failed` (terminal non-`COMPLETED`), carrying the problem's `detail` as `message`, its `run_status` member as `status` and its `error` member, the run's stored report, typed + Args: + run_id: The run to read. + artifacts: Which result artifacts to read, sent as one comma-separated `?artifacts=` + parameter. `None` (the default) reads them all. With a selection the platform + reads, re-signs and returns only those: an artifact left out is absent from the + result (`RunResults.carries` answers `False`), one asked for but never written is + `None`. A history row that shows a run's output asks for `[RunArtifact.MAIN_STUFF]` + alone instead of paying for the graph and the forms. The main-stuff check below + applies only when `main_stuff` was asked for. + Raises: + PipelineRequestError: If `artifacts` is an empty selection, which names nothing to read. + MissingMainStuffError: If a completed run asked for its main stuff delivers none. RunLifecycleUnavailableError: If the lifecycle routes are absent (a bare runner). ApiUnreachableError: If the host cannot be reached (DNS / connect / TLS / timeout). - ApiResponseError: For a genuine run-not-found 404 or any other non-2xx response. + ApiResponseError: For a genuine run-not-found 404 or any other non-2xx response + (an artifact name the platform does not know is its `400`). """ + selection = _artifact_selection(artifacts) endpoint = f"{_RUNS}/{quote(run_id, safe='')}/results" + if selection is not None: + endpoint = f"{endpoint}?{urlencode({'artifacts': ','.join(selection)}, safe=',')}" url = self._url(endpoint) response = await self._send_or_unreachable("GET", url, content=None, request_timeout=_POLL_REQUEST_TIMEOUT_SECONDS) status_code = response.status_code @@ -827,18 +844,25 @@ async def get_run_result(self, run_id: str) -> RunResultState: self._raise_if_lifecycle_unavailable(status=response.status_code, body=response.text, url=url) if not response.is_success: self._raise_api_response_error(method="GET", endpoint=endpoint, response=response) - # Inspect the decoded payload before validating: `main_stuff` is a required field on `RunResults`, - # so a `200` that omits the key would raise a raw Pydantic error instead of the typed - # `MissingMainStuffError`. `.get(...) is None` covers both the missing-key and explicit-null cases - # for the same un-deliverable-output condition (a present-but-falsy main stuff — `[]`, `0` — stays). + # A completed run asked for its main stuff must deliver one. `.get(...) is None` covers both the + # missing-key and explicit-null cases for the same un-deliverable-output condition (a + # present-but-falsy main stuff — `[]`, `0` — stays). A selection that left `main_stuff` out + # asked for none, so its absence there is the answer, not a fault. payload = response.json() - if isinstance(payload, dict) and cast("dict[str, Any]", payload).get("main_stuff") is None: + wants_main_stuff = selection is None or RunArtifact.MAIN_STUFF in selection + if wants_main_stuff and isinstance(payload, dict) and cast("dict[str, Any]", payload).get("main_stuff") is None: msg = f"Completed run '{run_id}' returned no main stuff — a completed run always delivers a main stuff." raise MissingMainStuffError(msg, run_id=run_id) result = RunResults.model_validate(payload) return RunResultCompleted(pipeline_run_id=run_id, result=result) - async def wait_for_result(self, run_id: str, options: WaitForResultOptions | None = None) -> RunResults: + async def wait_for_result( + self, + run_id: str, + options: WaitForResultOptions | None = None, + *, + artifacts: Sequence[RunArtifact] | None = None, + ) -> RunResults: """Poll a run to a terminal state and return its result. Resolves on `COMPLETED`, raises `RunFailedError` on any other terminal status — carrying the @@ -846,7 +870,10 @@ async def wait_for_result(self, run_id: str, options: WaitForResultOptions | Non `RunTimeoutError` if `timeout_seconds` elapses first (the run keeps executing server-side — resume later by `run_id`). Honors the server's `Retry-After`. Async-native: cancelling the awaiting task raises `asyncio.CancelledError` out of this loop, leaving the run resumable. + `artifacts` narrows every results read of the loop, exactly as on `get_run_result`. """ + # Refused before the first poll, so an empty selection never waits out a timeout to fail. + _artifact_selection(artifacts) opts = options or WaitForResultOptions() started_at = monotonic() attempt = 0 @@ -858,7 +885,7 @@ async def wait_for_result(self, run_id: str, options: WaitForResultOptions | Non raise RunTimeoutError(_timeout_message(run_id, opts.timeout_seconds), run_id=run_id, timeout_seconds=opts.timeout_seconds) try: - state = await asyncio.wait_for(self.get_run_result(run_id), timeout=remaining) + state = await asyncio.wait_for(self.get_run_result(run_id, artifacts=artifacts), timeout=remaining) except TimeoutError as exc: raise RunTimeoutError(_timeout_message(run_id, opts.timeout_seconds), run_id=run_id, timeout_seconds=opts.timeout_seconds) from exc @@ -911,6 +938,7 @@ async def start_and_wait( *, method_ref: str | None = None, method_id: str | None = None, + artifacts: Sequence[RunArtifact] | None = None, ) -> RunResults: """Start a run and wait for its result — the whole lifecycle in one call, self-healing across hosted and bare runners. @@ -930,10 +958,16 @@ async def start_and_wait( blocking fallback, and dropping `method_id` there would turn a server-side 422 that names the key into a silently different run. + `artifacts` narrows the hosted results read, as on `get_run_result`. The blocking path + ignores it: the execute response already holds every artifact, so there is nothing to save + by narrowing it, and the result it returns answers for every field. + Raises: RunFailedError: If the run reaches a terminal status other than COMPLETED. RunTimeoutError: If the poll budget elapses (the run keeps executing — resume by id). """ + # Refused before anything starts, so an empty selection never costs a run. + _artifact_selection(artifacts) if await self._supports_run_lifecycle(): try: started = await self.start( @@ -960,7 +994,7 @@ async def start_and_wait( method_ref=method_ref, method_id=method_id, ) - return await self.wait_for_result(started.pipeline_run_id, options=wait_options) + return await self.wait_for_result(started.pipeline_run_id, options=wait_options, artifacts=artifacts) return await self._execute_blocking( pipe_code=pipe_code, @@ -1353,7 +1387,9 @@ async def list_runs( cursor: The `next_cursor` of the previous page, passed back opaquely. Returns: - A `RunPage` of `PipelineRun` rows. For the whole history, prefer `iterate_runs`. + A `RunPage` of `RunHistoryItem` rows — what a history row shows, nothing more. Open + one run with `get_run_detail` for the whole record. For the whole history, prefer + `iterate_runs`. Raises: ApiResponseError: On any non-2xx. Note that every `/v1/runs*` product route sits @@ -1371,7 +1407,7 @@ async def iterate_runs( created_from: str | None = None, created_to: str | None = None, limit: int | None = None, - ) -> AsyncIterator[PipelineRun]: + ) -> AsyncIterator[RunHistoryItem]: """Yield every run of a method, following the cursors — `GET /v1/runs`. The same loop as `iterate_methods` with one deliberate difference: an **empty page ends @@ -1571,6 +1607,23 @@ def _product_query(params: dict[str, str | int | None]) -> str: return "?" + urlencode(kept) +def _artifact_selection(artifacts: Sequence[RunArtifact] | None) -> tuple[RunArtifact, ...] | None: + """Normalise a results-read selection: `None` reads everything, anything else is deduplicated + into the enum's declaration order, so the same selection always builds the same query. + + An empty selection names nothing to read; the platform refuses it with a `400`, and it is + refused here first, before any request, with the request-shape error the client already raises + for an empty `validate_files`. + """ + if artifacts is None: + return None + requested = set(artifacts) + if not requested: + msg = "An artifact selection must name at least one RunArtifact; pass artifacts=None to read them all." + raise PipelineRequestError(msg) + return tuple(artifact for artifact in RunArtifact if artifact in requested) + + def _is_gateway_timeout(exc: ApiResponseError | httpx.TimeoutException, elapsed_seconds: float) -> bool: """Whether a failed blocking `execute` is the hosted gateway's ~30s synchronous cut-off. diff --git a/pipelex_sdk/product_models.py b/pipelex_sdk/product_models.py index 21b5b29..c37d286 100644 --- a/pipelex_sdk/product_models.py +++ b/pipelex_sdk/product_models.py @@ -12,7 +12,7 @@ models name exactly what the routes accept. These are Pipelex-branded (the hosted product surface), so they live in this SDK, -not in `mthds`. `PipelineRun.status` reuses the run-lifecycle `RunStatus`. +not in `mthds`. `RunHistoryItem.status` and `PipelineRun.status` reuse the run-lifecycle `RunStatus`. """ from __future__ import annotations @@ -594,18 +594,33 @@ class UploadedFile(BaseModel): # detail read, and the admin-update route. -class PipeStatus(StrEnum): - """Per-pipe progress marker surfaced in a run's `pipe_statuses` map.""" +class RunHistoryItem(BaseModel): + """One row of a method's run list — `GET /v1/runs?method_id=…`. - SCHEDULED = "scheduled" - RUNNING = "running" - SUCCEEDED = "succeeded" - FAILED = "failed" - SKIPPED = "skipped" + What a run-history row shows, and nothing else: which run, its status, when it started and + finished, which pipe it ran and, for a failed run, why. The rest of a run record — its + organization, its creator, its method (the caller just named it), its workflow id, its result + prefix — is either already known to the caller or only meaningful once a run is opened, which + `get_run_detail` serves. + """ + + model_config = ConfigDict(extra="allow") + + pipeline_run_id: str + status: RunStatus + created_at: str + finished_at: str | None = None + pipe_code: str | None = None + """The pipe that ran, when it was named. A run that let the bundle's `main_pipe` decide + has none to report, so the platform serves this as null.""" + + error: LenientRunErrorReport = None + """The run's stored report, kept on the row so opening a failed run from history shows why + without another read. `None` on every run that did not fail.""" class PipelineRun(BaseModel): - """One run record in a method's run list — `GET /v1/runs?method_id=…`.""" + """A whole run record — the base of `RunDetail`, the shape `GET /v1/runs/{id}` serves.""" model_config = ConfigDict(extra="allow") @@ -624,7 +639,6 @@ class PipelineRun(BaseModel): status: RunStatus result_url: str | None = None error: LenientRunErrorReport = None - pipe_statuses: dict[str, PipeStatus] | None = None created_at: str finished_at: str | None = None @@ -651,7 +665,7 @@ class RunPage(BaseModel): model_config = ConfigDict(extra="allow") - items: list[PipelineRun] + items: list[RunHistoryItem] next_cursor: str | None = None diff --git a/pipelex_sdk/runs.py b/pipelex_sdk/runs.py index 7c67daf..6c19a95 100644 --- a/pipelex_sdk/runs.py +++ b/pipelex_sdk/runs.py @@ -31,7 +31,7 @@ Wire contract mirrors `pipelex-platform`: POST /v1/start -> RunResultStart (start, 202) GET /v1/runs/{pipeline_run_id}/status -> RunRead (status, self-healing) - GET /v1/runs/{pipeline_run_id}/results -> 202 / 200 / 409 (results) + GET /v1/runs/{pipeline_run_id}/results -> 202 / 200 / 409 (results; `?artifacts=` narrows the 200) """ from __future__ import annotations @@ -99,6 +99,44 @@ def is_success(self) -> bool: return False +# ── Result artifacts ──────────────────────────────────────────────── + + +class RunArtifact(StrEnum): + """One result artifact the results read can be asked for, named by the response field it fills. + + `get_run_result(run_id, artifacts=[...])` sends the selection as `?artifacts=a,b`, and the + platform then reads, re-signs and returns only those artifacts: a field that was not asked for + is ABSENT from the body, where a field that was asked for but never written is `null`. + `TOKENS_USAGES` names the usage envelope, so asking for it fills `usage_assembly_error` too. + Mirrors the platform's own `RunArtifact`. + """ + + GRAPH_SPEC = "graph_spec" + PIPE_IO_CONTRACTS = "pipe_io_contracts" + INPUT_FORM = "input_form" + OUTPUT_FORM = "output_form" + MAIN_STUFF = "main_stuff" + WORKING_MEMORY = "working_memory" + TOKENS_USAGES = "tokens_usages" + + @property + def results_fields(self) -> tuple[str, ...]: + """The `RunResults` fields this artifact fills on the hosted results read.""" + match self: + case RunArtifact.TOKENS_USAGES: + return ("tokens_usages", "usage_assembly_error") + case ( + RunArtifact.GRAPH_SPEC + | RunArtifact.PIPE_IO_CONTRACTS + | RunArtifact.INPUT_FORM + | RunArtifact.OUTPUT_FORM + | RunArtifact.MAIN_STUFF + | RunArtifact.WORKING_MEMORY + ): + return (self,) + + # ── Responses ─────────────────────────────────────────────────────── @@ -221,34 +259,36 @@ class TokensUsageRecord(BaseModel): class RunResults(BaseModel): """Result artifacts for a completed run — `GET /v1/runs/{pipeline_run_id}/results`. - `main_stuff` is the resolved main output content and is ALWAYS present for a - completed run (the pipelex >= 0.37 main-stuff invariant): on the hosted path - it is the `main_stuff.json` S3 artifact relayed verbatim; on the bare-runner - blocking path the SDK resolves it from the returned working memory via the - run's `main_stuff_name`, so both paths deliver the same content shape. - Consumers read `main_stuff` directly — no shape-guessing. A completed run that - cannot deliver a main stuff raises `MissingMainStuffError`. Extension-open - (`extra="allow"`): any other server artifact the SDK does not name is preserved - on `model_extra` rather than dropped. - - Every field but the first two is optional, and two readings of an optional field - are distinct on purpose. A key the hosted body did not carry is not in - `model_fields_set` and reads `None`; a key relayed as `null` is in the set and - reads `None` too. That is how a reader tells "the platform relayed no such key" - from "the platform relayed null", where the JS twin reads `undefined` against - `null`. The blocking path always answers for every field, so each is set there. + `main_stuff` is the resolved main output content, and a completed run read in full ALWAYS + carries one (the pipelex >= 0.37 main-stuff invariant): on the hosted path it is the + `main_stuff.json` S3 artifact relayed verbatim; on the bare-runner blocking path the SDK resolves + it from the returned working memory via the run's `main_stuff_name`, so both paths deliver the + same content shape. A completed run that cannot deliver a main stuff it was asked for raises + `MissingMainStuffError`. Extension-open (`extra="allow"`): any other server artifact the SDK + does not name is preserved on `model_extra` rather than dropped. + + Every field but `pipeline_run_id` is optional, and two readings of an optional field are + distinct on purpose. A key the hosted body did not carry is not in `model_fields_set` and reads + `None`; a key relayed as `null` is in the set and reads `None` too. With an artifact selection + (`get_run_result(run_id, artifacts=[...])`) that is exactly "not requested" against "requested + but never written": the platform leaves an unselected artifact out of the body and relays a + selected-but-unwritten one as `null`. `carries(artifact)` asks the question by artifact rather + than by field name. The blocking path always answers for every field, so each is set there. Every field is walked on `docs/run-results.md`. """ model_config = ConfigDict(extra="allow") pipeline_run_id: str - #: The resolved main output content — always present for a completed run. Typed `Any` because the - #: content is polymorphic (a structured output is an object of the concept's fields, a multiple - #: output the `{"items": [...]}` envelope the runtime's `ListContent` serialises to, a native is - #: wrapped too — `{"text": ...}`, `{"number": ...}`) and may be a valid empty value (an empty - #: `items`, an empty `text`); it is never absent for a completed run. - main_stuff: Any + #: The resolved main output content — never `None` on a completed run whose read asked for it (no + #: selection, or one naming `RunArtifact.MAIN_STUFF`): the client raises `MissingMainStuffError` + #: rather than hand back a result without it. Absent (not in `model_fields_set`, reading `None`) + #: only when a selection left it out. Typed `Any` because the content is polymorphic (a + #: structured output is an object of the concept's fields, a multiple output the + #: `{"items": [...]}` envelope the runtime's `ListContent` serialises to, a native is wrapped + #: too — `{"text": ...}`, `{"number": ...}`) and may be a valid empty value (an empty `items`, + #: an empty `text`). + main_stuff: Any = None #: The executed graph — the same document a local run writes as `graphspec.json`: `meta.mode` #: `"live"`, one node per pipe with its status, its timings and its own usage. It reaches the #: client on both paths: the hosted path relays the `graphspec.json` artifact verbatim, and on @@ -325,6 +365,16 @@ class RunResults(BaseModel): #: `tokens_usages` as `None`, so a caller that cares must branch on this, not on the list. usage_assembly_error: str | None = None + def carries(self, artifact: RunArtifact) -> bool: + """Whether the results body carried `artifact` at all, whatever its value. + + `False` means the read did not ask for it (a selection that left it out), so its field reads + `None` without saying anything about the run. `True` with a `None` value means the artifact + was asked for and the platform relayed `null`: it was never written. Always `True` on the + blocking path, which answers for every field. + """ + return all(field_name in self.model_fields_set for field_name in artifact.results_fields) + # ── Single-shot result lookup outcome (discriminated on `state`) ───── diff --git a/tests/unit/test_artifacts.py b/tests/unit/test_artifacts.py index 94ede1c..37f0b52 100644 --- a/tests/unit/test_artifacts.py +++ b/tests/unit/test_artifacts.py @@ -49,10 +49,10 @@ RunStillRunningError, ScopeUnavailableError, ) -from pipelex_sdk.runs import RunResultCompleted, RunResultFailed, RunResultRunning, RunResults, RunStatus +from pipelex_sdk.runs import RunArtifact, RunResultCompleted, RunResultFailed, RunResultRunning, RunResults, RunStatus if TYPE_CHECKING: - from collections.abc import AsyncIterator, Callable + from collections.abc import AsyncIterator, Callable, Sequence from pathlib import Path from pytest_mock import MockerFixture @@ -88,6 +88,7 @@ def __init__( self._run_result = run_result self.resolve_calls: list[list[str]] = [] self.run_result_calls: list[str] = [] + self.run_result_selections: list[list[RunArtifact] | None] = [] async def resolve_storage_urls_bulk(self, uris: list[str]) -> BulkResolvedStorageUrls: self.resolve_calls.append(list(uris)) @@ -96,8 +97,9 @@ async def resolve_storage_urls_bulk(self, uris: list[str]) -> BulkResolvedStorag raise AssertionError(msg) return self._resolve(list(uris)) - async def get_run_result(self, run_id: str) -> RunResultState: + async def get_run_result(self, run_id: str, *, artifacts: Sequence[RunArtifact] | None = None) -> RunResultState: self.run_result_calls.append(run_id) + self.run_result_selections.append(None if artifacts is None else list(artifacts)) if self._run_result is None: msg = "this test did not script a run result" raise AssertionError(msg) @@ -745,6 +747,7 @@ def _handler(request: httpx.Request) -> httpx.Response: verdict = asyncio.run(download_artifacts(client, dir_path=target, run_id=_RUN_ID)) assert client.run_result_calls == [_RUN_ID] + assert client.run_result_selections == [[RunArtifact.MAIN_STUFF]] assert client.resolve_calls == [[_URI_PNG, _URI_PDF]] assert verdict.all_saved is True assert [artifact.uri for artifact in verdict.artifacts] == [_URI_PNG, _URI_PDF] @@ -755,6 +758,23 @@ def _handler(request: httpx.Request) -> httpx.Response: assert (target / "items-0.png").read_bytes() == _PNG_BYTES assert (target / "items-1.pdf").read_bytes() == _PDF_BYTES + def test_reads_by_run_id_asks_only_for_the_working_memory_it_walks(self, mocker: MockerFixture, tmp_path: Path) -> None: + """A working-memory download reads that one artifact, so a result without a main stuff is no fault.""" + working_memory = {"root": {"doc": {"concept": "native.Document", "content": _content(_URI_PDF)}}, "aliases": {}} + results = RunResults.model_validate({"pipeline_run_id": _RUN_ID, "working_memory": working_memory}) + client = _FakeClient( + resolve=_resolver(_resolved(_URI_PDF)), + run_result=RunResultCompleted(pipeline_run_id=_RUN_ID, result=results), + ) + _patch_storage(mocker, _serving(_PDF_BYTES)) + options = DownloadArtifactsOptions(scope=ArtifactScope.WORKING_MEMORY) + verdict = asyncio.run(download_artifacts(client, dir_path=tmp_path / "out", run_id=_RUN_ID, options=options)) + + assert client.run_result_selections == [[RunArtifact.WORKING_MEMORY]] + assert verdict.scope == ArtifactScope.WORKING_MEMORY + assert verdict.all_saved is True + assert [artifact.uri for artifact in verdict.artifacts] == [_URI_PDF] + def test_takes_results_in_hand_without_re_reading_and_creates_the_directory(self, mocker: MockerFixture, tmp_path: Path) -> None: client = _FakeClient(resolve=_resolver(_resolved(_URI_PDF))) _patch_storage(mocker, _serving(_PDF_BYTES)) diff --git a/tests/unit/test_client_lifecycle.py b/tests/unit/test_client_lifecycle.py index 4b4fed9..81694e0 100644 --- a/tests/unit/test_client_lifecycle.py +++ b/tests/unit/test_client_lifecycle.py @@ -19,6 +19,7 @@ ) from pipelex_sdk.runs import ( PollInfo, + RunArtifact, RunResultCompleted, RunResultFailed, RunResultRunning, @@ -336,6 +337,108 @@ def test_get_run_result_completed_keeps_falsy_main_stuff(self, mocker: MockerFix assert isinstance(state, RunResultCompleted) assert state.result.main_stuff == [] + # ── get_run_result artifact selection ──────────────────────── + + @pytest.mark.parametrize( + ("artifacts", "expected_query"), + [ + pytest.param([RunArtifact.MAIN_STUFF], "artifacts=main_stuff", id="one"), + pytest.param( + [RunArtifact.TOKENS_USAGES, RunArtifact.MAIN_STUFF, RunArtifact.GRAPH_SPEC], + "artifacts=graph_spec,main_stuff,tokens_usages", + id="declaration_order", + ), + pytest.param([RunArtifact.MAIN_STUFF, RunArtifact.MAIN_STUFF], "artifacts=main_stuff", id="deduplicated"), + pytest.param(list(RunArtifact), "artifacts=" + ",".join(RunArtifact), id="all_named"), + ], + ) + def test_get_run_result_sends_the_selection_as_one_comma_separated_param( + self, mocker: MockerFixture, artifacts: list[RunArtifact], expected_query: str + ) -> None: + """A selection is one `artifacts` param, in declaration order, each name once, commas unescaped.""" + client = self._client() + body: dict[str, object] = {"pipeline_run_id": "run_1", "main_stuff": {"text": "hi"}} + send_mock = mocker.patch.object(client, "_send", mocker.AsyncMock(return_value=_response(200, json=body))) + + asyncio.run(client.get_run_result("run_1", artifacts=artifacts)) + assert send_mock.call_args.args[1] == f"{_BASE_URL}/v1/runs/run_1/results?{expected_query}" + + def test_get_run_result_without_a_selection_sends_no_query(self, mocker: MockerFixture) -> None: + """No selection reads everything, exactly as before: the URL carries no `artifacts` param.""" + client = self._client() + body: dict[str, object] = {"pipeline_run_id": "run_1", "main_stuff": {"text": "hi"}} + send_mock = mocker.patch.object(client, "_send", mocker.AsyncMock(return_value=_response(200, json=body))) + + state = asyncio.run(client.get_run_result("run_1")) + assert send_mock.call_args.args[1] == f"{_BASE_URL}/v1/runs/run_1/results" + assert isinstance(state, RunResultCompleted) + assert state.result.carries(RunArtifact.MAIN_STUFF) is True + + def test_get_run_result_refuses_an_empty_selection_before_sending(self, mocker: MockerFixture) -> None: + """An empty selection names nothing; it is refused client-side and no request is made.""" + client = self._client() + send_mock = mocker.patch.object(client, "_send", mocker.AsyncMock()) + + with pytest.raises(PipelineRequestError, match="at least one RunArtifact"): + asyncio.run(client.get_run_result("run_1", artifacts=[])) + send_mock.assert_not_called() + + def test_get_run_result_selection_without_main_stuff_does_not_require_one(self, mocker: MockerFixture) -> None: + """A selection that leaves `main_stuff` out reads a body without it, and that is the answer, not a fault.""" + client = self._client() + body: dict[str, object] = {"pipeline_run_id": "run_1", "graph_spec": {"nodes": []}, "output_form": None} + mocker.patch.object(client, "_send", mocker.AsyncMock(return_value=_response(200, json=body))) + + state = asyncio.run(client.get_run_result("run_1", artifacts=[RunArtifact.GRAPH_SPEC, RunArtifact.OUTPUT_FORM])) + assert isinstance(state, RunResultCompleted) + result = state.result + assert result.graph_spec == {"nodes": []} + assert result.main_stuff is None + assert result.carries(RunArtifact.GRAPH_SPEC) is True + assert result.carries(RunArtifact.OUTPUT_FORM) is True + assert result.output_form is None + assert result.carries(RunArtifact.MAIN_STUFF) is False + assert result.carries(RunArtifact.WORKING_MEMORY) is False + + @pytest.mark.parametrize( + "body", + [ + pytest.param({"pipeline_run_id": "run_1"}, id="main_stuff_key_omitted"), + pytest.param({"pipeline_run_id": "run_1", "main_stuff": None}, id="main_stuff_null"), + ], + ) + def test_get_run_result_selection_naming_main_stuff_still_requires_one(self, mocker: MockerFixture, body: dict[str, object]) -> None: + """A selection that asks for `main_stuff` keeps the invariant: none delivered is `MissingMainStuffError`.""" + client = self._client() + mocker.patch.object(client, "_send", mocker.AsyncMock(return_value=_response(200, json=body))) + + with pytest.raises(MissingMainStuffError) as exc_info: + asyncio.run(client.get_run_result("run_1", artifacts=[RunArtifact.MAIN_STUFF, RunArtifact.TOKENS_USAGES])) + assert exc_info.value.run_id == "run_1" + + def test_get_run_result_carries_the_usage_pair_only_when_both_fields_came(self, mocker: MockerFixture) -> None: + """`TOKENS_USAGES` names the envelope: it is carried when the body holds both of its fields.""" + client = self._client() + body: dict[str, object] = {"pipeline_run_id": "run_1", "main_stuff": {"text": "hi"}, "tokens_usages": None, "usage_assembly_error": None} + mocker.patch.object(client, "_send", mocker.AsyncMock(return_value=_response(200, json=body))) + + state = asyncio.run(client.get_run_result("run_1", artifacts=[RunArtifact.MAIN_STUFF, RunArtifact.TOKENS_USAGES])) + assert isinstance(state, RunResultCompleted) + assert state.result.carries(RunArtifact.TOKENS_USAGES) is True + assert state.result.tokens_usages is None + partial = RunResults.model_validate({"pipeline_run_id": "run_1", "tokens_usages": []}) + assert partial.carries(RunArtifact.TOKENS_USAGES) is False + + def test_get_run_result_with_a_selection_maps_a_202_to_running(self, mocker: MockerFixture) -> None: + """The in-flight body names the selected fields as null; the client still reads it as running.""" + client = self._client() + body: dict[str, object] = {"pipeline_run_id": "run_1", "main_stuff": None} + mocker.patch.object(client, "_send", mocker.AsyncMock(return_value=_response(202, json=body, headers={"Retry-After": "4"}))) + + state = asyncio.run(client.get_run_result("run_1", artifacts=[RunArtifact.MAIN_STUFF])) + assert isinstance(state, RunResultRunning) + assert state.retry_after_seconds == 4 + def test_get_run_result_completed_parses_usage_pair(self, mocker: MockerFixture) -> None: """A 200 carrying the hosted usage pair validates the relayed records into `TokensUsageRecord`s; a body without them (older platform / pre-artifact run) defaults both @@ -552,6 +655,50 @@ def test_wait_for_result_polls_until_completed(self, mocker: MockerFixture) -> N assert len(polls) == 1 assert polls[0].attempt == 1 + def test_wait_for_result_threads_the_selection_to_every_poll(self, mocker: MockerFixture) -> None: + """Every results read of the loop carries the caller's selection.""" + client = self._client() + result = RunResults(pipeline_run_id="run_1", main_stuff={"answer": "42"}) + get_mock = mocker.patch.object( + client, + "get_run_result", + mocker.AsyncMock( + side_effect=[ + RunResultRunning(pipeline_run_id="run_1", retry_after_seconds=0), + RunResultCompleted(pipeline_run_id="run_1", result=result), + ] + ), + ) + mocker.patch("pipelex_sdk.client.asyncio.sleep", mocker.AsyncMock()) + + asyncio.run(client.wait_for_result("run_1", WaitForResultOptions(interval_seconds=0.0), artifacts=[RunArtifact.MAIN_STUFF])) + assert [call.kwargs["artifacts"] for call in get_mock.call_args_list] == [[RunArtifact.MAIN_STUFF], [RunArtifact.MAIN_STUFF]] + + def test_wait_for_result_refuses_an_empty_selection_before_polling(self, mocker: MockerFixture) -> None: + """An empty selection fails at once rather than after a poll.""" + client = self._client() + get_mock = mocker.patch.object(client, "get_run_result", mocker.AsyncMock()) + + with pytest.raises(PipelineRequestError): + asyncio.run(client.wait_for_result("run_1", artifacts=[])) + get_mock.assert_not_called() + + def test_start_and_wait_threads_the_selection_to_the_wait(self, mocker: MockerFixture) -> None: + """On the hosted path the selection reaches `wait_for_result`; an empty one never starts a run.""" + client = self._client() + mocker.patch.object(client, "_supports_run_lifecycle", mocker.AsyncMock(return_value=True)) + start_mock = mocker.patch.object(client, "start", mocker.AsyncMock(return_value=mocker.Mock(pipeline_run_id="run_1"))) + result = RunResults(pipeline_run_id="run_1", main_stuff={"answer": "42"}) + wait_mock = mocker.patch.object(client, "wait_for_result", mocker.AsyncMock(return_value=result)) + + returned = asyncio.run(client.start_and_wait(pipe_code="p", artifacts=[RunArtifact.MAIN_STUFF])) + assert returned is result + assert wait_mock.call_args.kwargs["artifacts"] == [RunArtifact.MAIN_STUFF] + + with pytest.raises(PipelineRequestError): + asyncio.run(client.start_and_wait(pipe_code="p", artifacts=[])) + assert start_mock.call_count == 1 + def test_wait_for_result_raises_run_failed(self, mocker: MockerFixture) -> None: """A terminal non-COMPLETED state raises RunFailedError carrying the typed status.""" client = self._client() diff --git a/tests/unit/test_client_paging.py b/tests/unit/test_client_paging.py index 9955141..2ba0326 100644 --- a/tests/unit/test_client_paging.py +++ b/tests/unit/test_client_paging.py @@ -21,7 +21,7 @@ from pytest_mock import MockerFixture, MockType from pipelex_sdk.client import PipelexAPIClient - from pipelex_sdk.product_models import MethodSummary, PipelineRun + from pipelex_sdk.product_models import MethodSummary, RunHistoryItem from tests.unit.conftest import ResponseBuilder, SendPatcher @@ -30,7 +30,7 @@ def _method(method_id: str) -> dict[str, Any]: def _run(run_id: str) -> dict[str, Any]: - return {"pipeline_run_id": run_id, "method_id": "m1", "pipe_code": "p", "status": "RUNNING", "created_at": "t"} + return {"pipeline_run_id": run_id, "pipe_code": "p", "status": "RUNNING", "created_at": "t"} def _cursors_sent(send: MockType) -> list[str | None]: @@ -48,7 +48,7 @@ async def _drain_methods(client: PipelexAPIClient, **kwargs: Any) -> list[Method return [summary async for summary in client.iterate_methods(**kwargs)] -async def _drain_runs(client: PipelexAPIClient, method_id: str) -> list[PipelineRun]: +async def _drain_runs(client: PipelexAPIClient, method_id: str) -> list[RunHistoryItem]: return [pipeline_run async for pipeline_run in client.iterate_runs(method_id)] diff --git a/tests/unit/test_client_product.py b/tests/unit/test_client_product.py index 9b4a19e..21a26ca 100644 --- a/tests/unit/test_client_product.py +++ b/tests/unit/test_client_product.py @@ -24,6 +24,7 @@ OnboardingRole, OnboardingSubmission, OrgRole, + RunHistoryItem, UpdateRunInput, UploadInput, ) @@ -459,15 +460,15 @@ def test_list_runs_keeps_date_bounds_and_paging_params_on_presence(self, mocker: f"{_BASE_URL}/v1/runs?method_id=m1&created_from=2026-08-01T00%3A00%3A00%2B00%3A00&created_to=&limit=10&cursor=c1" ) - def test_run_row_parses_with_null_method_id_and_pipe_code(self, mocker: MockerFixture) -> None: - """An ad-hoc run belongs to no stored method, and a `main_pipe` run names no pipe.""" + def test_run_row_parses_the_history_fields_with_a_null_pipe_code(self, mocker: MockerFixture) -> None: + """A history row is exactly the history fields; a `main_pipe` run names no pipe.""" client = self._client() row = { "pipeline_run_id": "r1", - "method_id": None, "pipe_code": None, "status": "FAILED", - "created_at": "t", + "created_at": "2026-10-01T10:00:00Z", + "finished_at": "2026-10-01T10:00:05Z", "error": {"message": "boom", "error_type": "PipeExecutionError"}, } self._mock_send(mocker, client, _response(200, json_body={"items": [row], "next_cursor": None})) @@ -475,8 +476,14 @@ def test_run_row_parses_with_null_method_id_and_pipe_code(self, mocker: MockerFi result = asyncio.run(client.list_runs("m1")) pipeline_run = result.items[0] - assert pipeline_run.method_id is None + assert isinstance(pipeline_run, RunHistoryItem) + assert set(RunHistoryItem.model_fields) == {"pipeline_run_id", "status", "created_at", "finished_at", "pipe_code", "error"} + assert pipeline_run.pipeline_run_id == "r1" + assert pipeline_run.status == RunStatus.FAILED + assert pipeline_run.created_at == "2026-10-01T10:00:00Z" + assert pipeline_run.finished_at == "2026-10-01T10:00:05Z" assert pipeline_run.pipe_code is None + assert pipeline_run.model_extra == {} assert pipeline_run.error is not None assert pipeline_run.error.message == "boom" assert pipeline_run.error.error_type == "PipeExecutionError"