From 4b5fb1abe5dbf1e838e5d8488f452a6ddd6156f9 Mon Sep 17 00:00:00 2001 From: Louis Choquel Date: Tue, 22 Sep 2026 11:03:18 +0200 Subject: [PATCH] feature: the execute-to-RunResults lift becomes public `_map_run_result_to_run_results` was a private function in `client.py` with one caller, the blocking fallback inside `start_and_wait`. A consumer driving `execute()` itself and wanting `summarize_usage`, `collect_artifacts` or any run-results parity field had to re-read `pipe_output.model_extra` and re-validate the records by hand, which is what the Python starter does today. It moves to `pipelex_sdk/execute_result.py` as `results_from_execute(result)`, beside the `PipelexExecuteResult` it takes and on the side of the import edge that already depends on `runs`. The mapping itself is unchanged, the fallback calls the public name, and `execute()` still returns `PipelexExecuteResult` because that model carries the runner's whole typed envelope. Documented in the blocking-path opener of `docs/run-results.md`, in the bare-runner bullet of `docs/run-usage.md` and in `docs/architecture.md`; `tests/unit/test_results_from_execute.py` pins the function called on its own, including that its result feeds `summarize_usage`. Co-Authored-By: Claude Opus 5 (1M context) --- CHANGELOG.md | 6 + docs/architecture.md | 2 +- docs/run-results.md | 13 ++- docs/run-usage.md | 2 +- pipelex_sdk/client.py | 51 +-------- pipelex_sdk/execute_result.py | 66 ++++++++++- tests/unit/test_results_from_execute.py | 144 ++++++++++++++++++++++++ 7 files changed, 228 insertions(+), 56 deletions(-) create mode 100644 tests/unit/test_results_from_execute.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 8c4e029..5fe5cc3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,11 @@ # Changelog +## [Unreleased] + +### Added + +- **`results_from_execute`, the lift from a blocking result onto `RunResults`, made public.** `pipelex_sdk.execute_result.results_from_execute(result)` turns a `PipelexExecuteResult` into the `RunResults` the durable path hands back, lifting the usage pair, the graph pair, the working memory and the three I/O artifacts off the runner's extension-open `pipe_output` onto their declared fields. It is the same mapping `start_and_wait` has always applied on its bare-runner fallback, which until now was private: a caller driving the blocking `execute()` itself had to re-read `pipe_output.model_extra` and re-validate the records by hand to reach `summarize_usage`, `collect_artifacts` or any parity field. `execute()` still returns `PipelexExecuteResult`, since that model carries the runner's whole typed envelope. The function is pure — no client, no network — and is documented in `docs/run-results.md` and `docs/run-usage.md`. + ## [v0.10.1] - 2026-09-22 ### Added diff --git a/docs/architecture.md b/docs/architecture.md index 24203e9..c4f8e54 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -77,7 +77,7 @@ The protocol `execute` is **overridden** to add one Pipelex-API behavior the bar Everything else stays the inherited regime: the protocol's optional 202 async-degrade still raises `RunStillRunningError` (from the base `execute`), and every other non-2xx keeps `httpx.HTTPStatusError` — consistent with the other inherited protocol routes (decision #5). The `_execute_blocking` bare-runner fallback (below) calls this same overridden `execute`, so it inherits the translation; off-platform there is no gateway cap, so the `503/504`-after-28s condition effectively never fires there. -The override also **enriches the return type**: it re-validates the base result into a `PipelexExecuteResult` (`pipelex_sdk/execute_result.py`), a `DictRunResultExecute` subtype that adds a resolved `.main_stuff` accessor (dug out of `pipe_output`'s working memory via the response's `main_stuff_name`, raising `MissingMainStuffError` if unlocatable). This gives blocking and durable results the **same output accessor** — `result.main_stuff` — so callers never branch on which path ran. `_map_run_result_to_run_results` reads that accessor too, keeping the resolution single-sourced. +The override also **enriches the return type**: it re-validates the base result into a `PipelexExecuteResult` (`pipelex_sdk/execute_result.py`), a `DictRunResultExecute` subtype that adds a resolved `.main_stuff` accessor (dug out of `pipe_output`'s working memory via the response's `main_stuff_name`, raising `MissingMainStuffError` if unlocatable). This gives blocking and durable results the **same output accessor** — `result.main_stuff` — so callers never branch on which path ran. `results_from_execute` (the public lift from that result onto `RunResults`, in the same module) reads that accessor too, keeping the resolution single-sourced; the `_execute_blocking` fallback is its only caller inside the SDK, and a consumer driving `execute()` directly calls it by name. ## Method selectors on the run routes (`method_ref` + `method_id`) diff --git a/docs/run-results.md b/docs/run-results.md index 4be51ab..f652c16 100644 --- a/docs/run-results.md +++ b/docs/run-results.md @@ -4,6 +4,17 @@ A completed run hands back one object, `RunResults` (`pipelex_sdk/runs.py`), and Two paths produce it. Against the hosted API the SDK starts a durable run and polls `GET /v1/runs/{id}/results`, where the platform relays the run's S3 artifacts verbatim. Against a bare `pipelex-api` runner, which has no run store, the SDK falls back to the blocking `POST /v1/execute` and maps the runner's native `pipe_output` onto the same shape, lifting the artifacts that ride it onto their own fields. `start_and_wait` picks between the two from the `GET /v1/version` handshake, so a consumer does not choose. +**Calling the blocking route yourself.** `execute()` returns a `PipelexExecuteResult` rather than a `RunResults`, because that model is the runner's whole typed envelope and nothing of it is thrown away. To read such a result through this page's fields, lift it: `results_from_execute(result)` (`pipelex_sdk/execute_result.py`) is the same mapping `start_and_wait` applies on its fallback, exposed for the caller who drives `execute()` directly. It is pure — no client, no network — and what it buys is everything written against `RunResults`: `summarize_usage`, `download_artifacts`, the graph pair and the three I/O artifacts, instead of re-reading `pipe_output.model_extra` by hand. + +```python +from pipelex_sdk.execute_result import results_from_execute +from pipelex_sdk.usage import summarize_usage + +execute_result = await client.execute(pipe_code="my_domain.my_pipe", mthds_contents=[source]) +results = results_from_execute(execute_result) +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 | @@ -160,7 +171,7 @@ The usage pair reports what each inference call consumed and cost — one `Token ## `pipe_output` — the runner's native output -`pipe_output` is the bare runner's whole native output, and it is present on the blocking path only — the hosted results body carries no such key, so on that path it reads `None`. It is supplementary: `main_stuff`, the graph pair, the three I/O artifacts, the working memory and the usage pair are all lifted out of it onto fields that read the same on both paths, so a consumer that reads those fields keeps working against the hosted API. What `pipe_output` adds is the runner's output exactly as it arrived, typed as the standard's `DictPipeOutputAbstract`, which is extension-open — the runner's Pipelex extension fields, the `pipe_io_artifacts` envelope among them, stay reachable in their raw form through `model_extra`. +`pipe_output` is the bare runner's whole native output, and it is present on the blocking path only — the hosted results body carries no such key, so on that path it reads `None`. It is supplementary: `main_stuff`, the graph pair, the three I/O artifacts, the working memory and the usage pair are all lifted out of it onto fields that read the same on both paths, so a consumer that reads those fields keeps working against the hosted API. What `pipe_output` adds is the runner's output exactly as it arrived, typed as the standard's `DictPipeOutputAbstract`, which is extension-open — the runner's Pipelex extension fields, the `pipe_io_artifacts` envelope among them, stay reachable in their raw form through `model_extra`. A caller holding an `execute()` result rather than a `RunResults` does that lift with `results_from_execute`, described at the top of this page, instead of reading the bag itself. ## Produced files diff --git a/docs/run-usage.md b/docs/run-usage.md index eea5340..cfe5ac1 100644 --- a/docs/run-usage.md +++ b/docs/run-usage.md @@ -19,7 +19,7 @@ For the run's totals, do not add the records up by hand: [`summarize_usage`](#su The accessor is the same whichever path ran. `start_and_wait` picks a path from the `GET /v1/version` handshake: - **Hosted (durable) path** — the records come from the runner's `tokens_usages.json` artifact, which `GET /v1/runs/{id}/results` unpacks onto the results body as top-level keys and relays verbatim. -- **Bare runner (blocking) path** — the records ride the execute response's extension-open `pipe_output` as Pipelex extension fields; the SDK lifts them onto the same two top-level fields. +- **Bare runner (blocking) path** — the records ride the execute response's extension-open `pipe_output` as Pipelex extension fields; the SDK lifts them onto the same two top-level fields. A caller that drives the blocking `execute()` itself applies that same lift with `results_from_execute(result)` (`pipelex_sdk/execute_result.py`), which hands back the `RunResults` this page's accessors are written against — there is no need to read the records out of `model_extra` and validate them by hand. See [`run-results.md`](./run-results.md). Because the runtime emits both surfaces through one helper, the two cannot structurally diverge. diff --git a/pipelex_sdk/client.py b/pipelex_sdk/client.py index f076704..77cd024 100644 --- a/pipelex_sdk/client.py +++ b/pipelex_sdk/client.py @@ -62,7 +62,7 @@ RunLifecycleUnavailableError, RunTimeoutError, ) -from pipelex_sdk.execute_result import PipelexExecuteResult +from pipelex_sdk.execute_result import PipelexExecuteResult, results_from_execute from pipelex_sdk.prepare_inputs import PreparedInputs from pipelex_sdk.prepare_inputs import prepare_inputs as _prepare_inputs_impl from pipelex_sdk.product_models import ( @@ -937,7 +937,7 @@ async def _execute_blocking( method_ref=method_ref, method_id=method_id, ) - return _map_run_result_to_run_results(result) + return results_from_execute(result) # ── Pipelex product surface (hosted management routes) ───────────────── # @@ -1592,53 +1592,6 @@ def _extract_run_status_from_message(message: str) -> RunStatus: return RunStatus.FAILED -def _map_run_result_to_run_results(response: PipelexExecuteResult) -> RunResults: - """Map the protocol's blocking `POST /v1/execute` response onto the lifecycle's `RunResults`. - - `response.main_stuff` resolves the main output out of the returned working memory (and raises - `MissingMainStuffError` if the run named no locatable main stuff), so the durable and blocking - paths hand back the same `main_stuff` content shape. The already-parsed `pipe_output` model is - carried over as-is — no `.model_dump()` round-trip — so the runner's whole envelope stays typed - (blocking only; the hosted path has none), and its `working_memory` is lifted onto the field of - that name, where the hosted path relays the artifact as its own key. The standard declares - `DictPipeOutputAbstract.working_memory` required, so that lift always carries a value here. - - The graph pair (`graph_spec` / `graph_assembly_error`), the usage pair (`tokens_usages` / - `usage_assembly_error`), the `pipe_io_artifacts` envelope and its `pipe_io_artifacts_error` all - ride the execute response's extension-open `pipe_output` as Pipelex extension fields. Lifting - each onto its top-level field here is what makes `.graph_spec`, `.tokens_usages` and - `.pipe_io_contracts` read the same on the blocking and durable paths; `RunResults` validates the - raw records into `TokensUsageRecord`s and the raw artifacts into the standard's models on the - way in. The runner carries the three I/O artifacts in one envelope — they share a key set and - are always built together — where the hosted results body relays them as three sibling keys; - `RunResults` follows the hosted shape and this unwraps the envelope onto it. A null envelope - leaves all three `None`, beside whatever `pipe_io_artifacts_error` says about why. - - Every lifted field is passed explicitly, so on this path each is in `model_fields_set` whether - or not the runner carried the key — the blocking path always answers, as the JS twin writes - `null` there; only the hosted path leaves a field unset when the body did not carry its key. - """ - pipe_output_extras: dict[str, Any] = response.pipe_output.model_extra or {} - pipe_io_artifacts: dict[str, Any] = {} - raw_pipe_io_artifacts = pipe_output_extras.get("pipe_io_artifacts") - if isinstance(raw_pipe_io_artifacts, dict): - pipe_io_artifacts = cast("dict[str, Any]", raw_pipe_io_artifacts) - return RunResults( - pipeline_run_id=response.pipeline_run_id, - main_stuff=response.main_stuff, - graph_spec=pipe_output_extras.get("graph_spec"), - graph_assembly_error=pipe_output_extras.get("graph_assembly_error"), - pipe_io_contracts=pipe_io_artifacts.get("pipe_io_contracts"), - input_form=pipe_io_artifacts.get("input_form"), - output_form=pipe_io_artifacts.get("output_form"), - pipe_io_artifacts_error=pipe_output_extras.get("pipe_io_artifacts_error"), - working_memory=response.pipe_output.working_memory, - pipe_output=response.pipe_output, - tokens_usages=pipe_output_extras.get("tokens_usages"), - usage_assembly_error=pipe_output_extras.get("usage_assembly_error"), - ) - - def _is_valid_base_url(value: str) -> bool: """Whether a base URL is host-only — http/https, no path, query, fragment, or embedded credentials (auth travels in the Authorization header, never the URL). diff --git a/pipelex_sdk/execute_result.py b/pipelex_sdk/execute_result.py index 12db661..c103b99 100644 --- a/pipelex_sdk/execute_result.py +++ b/pipelex_sdk/execute_result.py @@ -1,17 +1,20 @@ -"""The blocking `execute()` result — a `DictRunResultExecute` that resolves its `.main_stuff`. +"""The blocking `execute()` result — a `DictRunResultExecute` that resolves its `.main_stuff`, and +the public lift from that result onto `RunResults`. Kept in its own module (not `runs.py`) so it can import `MissingMainStuffError` from -`errors` without forming an import cycle (`errors` type-imports `runs`). +`errors` without forming an import cycle (`errors` type-imports `runs`). The lift lives here for +the same reason and in the same direction: it takes a `PipelexExecuteResult` and builds a +`RunResults`, so it belongs on the side that already depends on `runs`. """ from __future__ import annotations -from typing import Any +from typing import Any, cast from mthds.runners.api.models import DictRunResultExecute from pipelex_sdk.errors import MissingMainStuffError -from pipelex_sdk.runs import MethodProvenance +from pipelex_sdk.runs import MethodProvenance, RunResults class PipelexExecuteResult(DictRunResultExecute): @@ -53,3 +56,58 @@ def main_stuff(self) -> Any: ) raise MissingMainStuffError(msg, run_id=self.pipeline_run_id) return stuff.content + + +def results_from_execute(result: PipelexExecuteResult) -> RunResults: + """Lift a blocking `execute()` result onto the lifecycle's `RunResults` — the same shape a durable run hands back. + + `execute()` returns the runner's whole typed envelope, where the usage pair, the graph pair, the + working memory and the three I/O artifacts ride the extension-open `pipe_output` as Pipelex + extension fields on `model_extra`. This function is the one place that lifts each onto the + declared field of the same name, so a caller driving the blocking route itself reaches + `pipelex_sdk.usage.summarize_usage`, `pipelex_sdk.artifacts.collect_artifacts` and every other + run-results field exactly as it would on the hosted path, instead of re-reading `model_extra` by + hand. It is pure: no client, no network, no I/O. `start_and_wait` calls it on its bare-runner + fallback, which is the only caller inside the SDK. + + `result.main_stuff` resolves the main output out of the returned working memory (and raises + `MissingMainStuffError` if the run named no locatable main stuff), so the durable and blocking + paths hand back the same `main_stuff` content shape. The already-parsed `pipe_output` model is + carried over as-is — no `.model_dump()` round-trip — so the runner's whole envelope stays typed + (blocking only; the hosted path has none), and its `working_memory` is lifted onto the field of + that name, where the hosted path relays the artifact as its own key. The standard declares + `DictPipeOutputAbstract.working_memory` required, so that lift always carries a value here. + + Lifting the graph pair (`graph_spec` / `graph_assembly_error`), the usage pair (`tokens_usages` / + `usage_assembly_error`) and the `pipe_io_artifacts` envelope with its `pipe_io_artifacts_error` + onto their top-level fields is what makes `.graph_spec`, `.tokens_usages` and `.pipe_io_contracts` + read the same on the blocking and durable paths; `RunResults` validates the raw records into + `TokensUsageRecord`s and the raw artifacts into the standard's models on the way in. The runner + carries the three I/O artifacts in one envelope — they share a key set and are always built + together — where the hosted results body relays them as three sibling keys; `RunResults` follows + the hosted shape and this unwraps the envelope onto it. A null envelope leaves all three `None`, + beside whatever `pipe_io_artifacts_error` says about why. + + Every lifted field is passed explicitly, so on this path each is in `model_fields_set` whether + or not the runner carried the key — the blocking path always answers, as the JS twin writes + `null` there; only the hosted path leaves a field unset when the body did not carry its key. + """ + pipe_output_extras: dict[str, Any] = result.pipe_output.model_extra or {} + pipe_io_artifacts: dict[str, Any] = {} + raw_pipe_io_artifacts = pipe_output_extras.get("pipe_io_artifacts") + if isinstance(raw_pipe_io_artifacts, dict): + pipe_io_artifacts = cast("dict[str, Any]", raw_pipe_io_artifacts) + return RunResults( + pipeline_run_id=result.pipeline_run_id, + main_stuff=result.main_stuff, + graph_spec=pipe_output_extras.get("graph_spec"), + graph_assembly_error=pipe_output_extras.get("graph_assembly_error"), + pipe_io_contracts=pipe_io_artifacts.get("pipe_io_contracts"), + input_form=pipe_io_artifacts.get("input_form"), + output_form=pipe_io_artifacts.get("output_form"), + pipe_io_artifacts_error=pipe_output_extras.get("pipe_io_artifacts_error"), + working_memory=result.pipe_output.working_memory, + pipe_output=result.pipe_output, + tokens_usages=pipe_output_extras.get("tokens_usages"), + usage_assembly_error=pipe_output_extras.get("usage_assembly_error"), + ) diff --git a/tests/unit/test_results_from_execute.py b/tests/unit/test_results_from_execute.py new file mode 100644 index 0000000..75effbf --- /dev/null +++ b/tests/unit/test_results_from_execute.py @@ -0,0 +1,144 @@ +"""Tests for `results_from_execute` — the public lift from a blocking execute result onto `RunResults`. + +The blocking fallback's own path is covered in `test_client_run_fallback.py`, through the client. What +this module pins is the function a direct `execute()` caller reaches for: called on its own, with no +client and no network, it must produce the same `RunResults` the durable path hands back, so +`summarize_usage` and the parity fields read the same whichever route ran. +""" + +from typing import Any + +import pytest + +from pipelex_sdk.errors import MissingMainStuffError +from pipelex_sdk.execute_result import PipelexExecuteResult, results_from_execute +from pipelex_sdk.usage import UsageSummaryState, summarize_usage + +_TOKENS_USAGES: list[dict[str, Any]] = [ + { + "model_type": "llm", + "inference_model_name": "test-model", + "inference_model_id": "test-model-2026-01-01", + "pipe_code": "test_domain.summarize", + "job_category": "llm_job", + "unit_job_id": "llm_gen_text", + "nb_tokens_by_category": {"input": 100, "output": 20}, + "cost": 0.01, + "started_at": "2026-09-22T10:00:01+00:00", + "completed_at": "2026-09-22T10:00:03+00:00", + } +] + +_PIPE_IO_ARTIFACTS: dict[str, Any] = { + "pipe_io_contracts": { + "x.greet": { + "inputs": {}, + "output": { + "concept_ref": "native.Text", + "multiplicity": "single", + "item_count": None, + "optional": False, + "json_schema": {"type": "object", "properties": {"text": {"type": "string"}}}, + }, + }, + }, + "input_form": {"x.greet": {"fields": []}}, + "output_form": {"x.greet": {"field": {"name": "text", "kind": "prose", "required": True}}}, +} + + +def _execute_result(**pipe_output_extension_fields: Any) -> PipelexExecuteResult: + """A completed blocking response, validated from the wire, with extension fields on `pipe_output`.""" + return PipelexExecuteResult.model_validate( + { + "pipeline_run_id": "run-x", + "main_stuff_name": "result", + "pipe_output": { + "pipeline_run_id": "run-x", + "working_memory": { + "root": {"result": {"concept": "native.Text", "content": {"text": "hello"}}}, + "aliases": {"main_stuff": "result"}, + }, + **pipe_output_extension_fields, + }, + } + ) + + +class TestResultsFromExecute: + def test_lifts_every_extension_field_onto_its_own_field(self) -> None: + """The usage pair, the graph pair and the `pipe_io_artifacts` envelope all come off `model_extra`.""" + results = results_from_execute( + _execute_result( + tokens_usages=_TOKENS_USAGES, + graph_spec={"meta": {"format": "mthds"}, "nodes": [], "edges": []}, + pipe_io_artifacts=_PIPE_IO_ARTIFACTS, + ) + ) + + assert results.pipeline_run_id == "run-x" + assert results.main_stuff == {"text": "hello"} + assert results.graph_spec == {"meta": {"format": "mthds"}, "nodes": [], "edges": []} + assert results.tokens_usages is not None + assert results.tokens_usages[0].cost == 0.01 + assert results.pipe_io_contracts is not None + assert "x.greet" in results.pipe_io_contracts + assert results.input_form is not None + assert results.output_form is not None + + def test_carries_the_working_memory_and_the_parsed_pipe_output(self) -> None: + """`working_memory` is lifted off the runner's output and `pipe_output` is carried over as-is.""" + result = _execute_result() + results = results_from_execute(result) + + assert results.working_memory is not None + assert results.working_memory.root["result"].content == {"text": "hello"} + assert results.pipe_output is result.pipe_output + + def test_every_lifted_field_is_set_even_when_the_runner_carried_none(self) -> None: + """The blocking path always answers: a field the runner did not carry is `None` AND in the set.""" + results = results_from_execute(_execute_result()) + + for field_name in ( + "graph_spec", + "graph_assembly_error", + "pipe_io_contracts", + "input_form", + "output_form", + "pipe_io_artifacts_error", + "tokens_usages", + "usage_assembly_error", + ): + assert getattr(results, field_name) is None + assert field_name in results.model_fields_set + + def test_a_null_artifacts_envelope_leaves_the_three_none_beside_its_error(self) -> None: + """A runner that failed to build the artifacts sends a null envelope and says why.""" + results = results_from_execute(_execute_result(pipe_io_artifacts=None, pipe_io_artifacts_error="build broke")) + + assert results.pipe_io_contracts is None + assert results.input_form is None + assert results.output_form is None + assert results.pipe_io_artifacts_error == "build broke" + + def test_the_lifted_results_feed_summarize_usage(self) -> None: + """The reason the lift is public: a direct `execute()` caller reaches the run-level usage reading.""" + summary = summarize_usage(results_from_execute(_execute_result(tokens_usages=_TOKENS_USAGES))) + + assert summary.state == UsageSummaryState.RECORDS + assert summary.total_cost_usd == 0.01 + assert summary.calls == 1 + assert summary.tokens.input == 100 + + def test_a_run_with_no_locatable_main_stuff_raises(self) -> None: + """`main_stuff` is resolved through `main_stuff_name`, so an unlocatable one fails here too.""" + result = PipelexExecuteResult.model_validate( + { + "pipeline_run_id": "run-x", + "main_stuff_name": "absent", + "pipe_output": {"pipeline_run_id": "run-x", "working_memory": {"root": {}, "aliases": {}}}, + } + ) + + with pytest.raises(MissingMainStuffError): + results_from_execute(result)