Skip to content
This repository was archived by the owner on Oct 2, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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, ...`
Expand Down
12 changes: 6 additions & 6 deletions docs/architecture.md

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion docs/artifact-download.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
24 changes: 21 additions & 3 deletions docs/run-results.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` |
Expand All @@ -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)
Expand All @@ -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.
Expand All @@ -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.

Expand Down
10 changes: 10 additions & 0 deletions pipelex_sdk/artifact_models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 ────────────────────────────────────────────────────────

Expand Down Expand Up @@ -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.
Expand Down
16 changes: 10 additions & 6 deletions pipelex_sdk/artifacts.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 ─────────────────────────────
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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 "."
Expand Down
Loading
Loading