Skip to content
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
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,17 @@
# Changelog

## [Unreleased]

### Changed

- **Pinned `pipelex` 0.68.0 (Breaking)**: up from `==0.67.0`, exactly, the release whose `execute` accepts the inbound request id the host is serving. The `.pipelex/` config shipped here already sits at the current schema, so no migration is required, but a deployment that selects the `gcp` log sink instead of `json` now refuses to boot when Google rejects its credentials, where it used to boot and lose every record in silence.
- **A log line's trace keys name the host's span, and the runtime's span moves under `pipelex.*` (Breaking)**: `trace_id`, `span_id` and `trace_flags` now name the process's current OpenTelemetry span, which this server does not open, so they are absent unless the deployment runs it under OpenTelemetry instrumentation of its own. A line the runtime emits from inside a traced run carries the runtime's span as `pipelex.trace_id` and `pipelex.span_id` instead, so a query that joined lines to the runtime's exported spans on `trace_id` now joins on `pipelex.trace_id`.
- **`POST /v1/codegen` stamps `engine_version` `0.68.0`**: the stamp is the pinned `pipelex` version, so a `codegen.lock` committed against `0.67.0` no longer matches until it is regenerated. `POST /v1/build/runner` carries the same stamp.

### Fixed

- **A distributed `POST /v1/execute` run's worker lines carry the request id**: the route now puts the request id it resolved on the run's metadata, as `POST /v1/start` already did, so on a deployment whose `orchestration_mode` dispatches to a worker, such as `temporal`, the lines the worker writes while it runs the run's workflow and activities carry the `request_id` the response echoes. The logging page now says how the id reaches a worker's lines, and that a `POST /v1/validate` dispatched to a worker does not carry it yet.

## [v0.29.0] - 2026-09-27

### Highlights
Expand Down
4 changes: 3 additions & 1 deletion api/middleware.py
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,9 @@ class RequestIdMiddleware:
puts `request_id` on *every* record emitted underneath — the API's own
error lines, and equally the ones pipelex emits from inside a run — as an
attribute a structured sink indexes, with no call site having to pass it
and no message having to interpolate it.
and no message having to interpolate it. The binding is in-process only: a
run dispatched to a worker gets the id from its `RunMetadata`, where the
run routes put it by passing `request_id_of(request)` to the runner.

Applied in `api.main` by wrapping the whole FastAPI app
(`app = RequestIdMiddleware(app)`), NOT via `app.add_middleware()`.
Expand Down
13 changes: 9 additions & 4 deletions api/routes/pipelex/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,7 @@ async def execute(
dynamic_output_concept_ref: str | None = None,
extra: dict[str, Any] | None = None,
delivery_assignment: DeliveryAssignment | None = None,
request_id: str | None = None,
requested_orchestration_mode: str | None = None,
) -> PipelexRunResultExecute:
"""Execute a method synchronously, dispatching by the resolved `orchestration_mode`.
Expand All @@ -238,8 +239,10 @@ async def execute(
The orchestrator is injected as this runner's `_pipe_run` so the inherited base `execute`
keeps the entire run lifecycle (library setup/teardown, tracer close, pipeline-manager
cleanup, telemetry, error mapping); only the dispatch backend and the output rehydration
(`_OrchestratorPipeRun`) change. `requested_orchestration_mode` is the optional per-request
backend override (`PipelineApiExtras.orchestration_mode`).
(`_OrchestratorPipeRun`) change. `request_id` is an API-layer extra threaded into
`RunMetadata.request_id` for log correlation, so it reaches the job a worker is handed.
`requested_orchestration_mode` is the optional per-request backend override
(`PipelineApiExtras.orchestration_mode`).
"""
# Resolve the effective orchestration mode FIRST — a per-request override the deployment
# policy forbids is refused (403) here, before any library load / run registration. Mirrors start().
Expand All @@ -262,6 +265,7 @@ async def execute(
dynamic_output_concept_ref=dynamic_output_concept_ref,
extra=extra,
delivery_assignment=delivery_assignment,
request_id=request_id,
)

@override
Expand Down Expand Up @@ -305,8 +309,8 @@ async def start(
`extra` is the protocol's generic extension slot; this runner's wire
extras are parsed by the route layer, so nothing reaches it — a
non-empty value is an in-process misuse and is rejected. `request_id`
is an API-layer extra threaded into `JobMetadata.request_id` for log
correlation. `requested_orchestration_mode` is the optional per-request backend override
is an API-layer extra threaded into `RunMetadata.request_id` for log
correlation, so it reaches the job a worker is handed. `requested_orchestration_mode` is the optional per-request backend override
(`PipelineApiExtras.orchestration_mode`); it is resolved against the deployment's
`api.toml` policy and a forbidden override is refused with a 403.
"""
Expand Down Expand Up @@ -747,6 +751,7 @@ async def execute(request: Request) -> JSONResponse:
output_name=run_request.output_name,
output_multiplicity=run_request.output_multiplicity,
dynamic_output_concept_ref=run_request.dynamic_output_concept_ref,
request_id=request_id_of(request),
requested_orchestration_mode=extras.orchestration_mode,
)
# The response dump carries the full internal usage models on
Expand Down
2 changes: 1 addition & 1 deletion docs/error-responses.md
Original file line number Diff line number Diff line change
Expand Up @@ -202,7 +202,7 @@ API-authored errors (`ValidationError`, `BadRequest`, `Unauthenticated`, etc.) f

## Request correlation

Every response carries `X-Request-ID`. The middleware respects an inbound `X-Request-ID` header if present, otherwise mints one. It also binds that id onto the Pipelex runtime's log context for the duration of the request, so every record emitted underneath carries it as a `request_id` field — the server's own error lines and the runtime's lines from inside a run alike. The same id rides onto `JobMetadata.request_id`, so on a distributed-execution flavor it correlates with every orchestrator worker-side record too. What those lines look like and which fields they carry is in [Logging](logging.md).
Every response carries `X-Request-ID`. The middleware respects an inbound `X-Request-ID` header if present, otherwise mints one. It also binds that id onto the Pipelex runtime's log context for the duration of the request, so every record emitted underneath carries it as a `request_id` field — the server's own error lines and the runtime's lines from inside a run alike. Both run routes, `POST /v1/execute` and `POST /v1/start`, also put the id on the run's metadata, so when a run is dispatched to a worker, the lines the worker writes while it runs the run's workflow and activities carry the same id. What those lines look like, which fields they carry and where the id does not reach are in [Logging](logging.md).

When opening an issue, include the `request_id` from the response (or response headers) and the timestamp.

Expand Down
7 changes: 5 additions & 2 deletions docs/logging.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,12 @@ These keys come from the sink itself and are on every line written by the proces
| `logger` | The module that emitted it |
| `message` | The human-readable summary |
| `exception` | The traceback, when the record carries one |
| `trace_id`, `span_id`, `trace_flags` | The trace context, in lowercase hex, when the record was logged inside a span, which is the case for a line the runtime emits from inside a traced run |
| `trace_id`, `span_id`, `trace_flags` | The process's current OpenTelemetry span, in lowercase hex, when one names a trace. This server registers no tracer and instruments nothing, so as shipped no line carries them; they appear only when the deployment runs it under OpenTelemetry instrumentation of its own that makes a span current |
| `pipelex.trace_id`, `pipelex.span_id` | The Pipelex span a runtime line was logged in, the pipe's or the LLM call's, in lowercase hex. Present on the runtime's lines inside a live run the runtime traces, which it does whenever the Pipelex Gateway's telemetry is on or AI span tracing is enabled on your own PostHog; never on the server's own lines, an error line included |

`request_id` comes from the request-scoped context the request-id middleware binds, so **every** record emitted while a request is in flight carries it. The example above is one of the server's own lines, but a line the Pipelex runtime emits from inside a pipeline run carries the same id, which is what ties the two together without any call site passing it along. The value is the one echoed in the response's `X-Request-ID` header and in the problem document's `request_id` member, so a caller reporting a failure hands you the key to its log lines.
The two pairs name different traces. Pipelex's spans belong to a trace of their own, which only Pipelex's exporters receive, so if you export them to a backend of your own, such as Langfuse, an OTLP collector or your PostHog, join a line to them on `pipelex.trace_id` and `pipelex.span_id`. To gather a run's lines whether or not anything traces, use `request_id` or `pipeline_run_id`. How each sink writes both pairs is in the Pipelex documentation's [trace context](https://docs.pipelex.com/0.68.0/tools/logging/#the-trace-context) section.

`request_id` reaches a line in one of two ways, depending on which process writes it. In this server's own process, the request-id middleware binds it on the runtime's log context for the whole request, so every record emitted there while the request is in flight carries it: the server's own lines, like the example above, and the lines the Pipelex runtime emits from inside a run executed in-process (`orchestration_mode = "direct"`), without any call site passing it along. A run dispatched to a worker, under a distributed `orchestration_mode` such as `temporal`, happens in another process, which that binding does not reach, so the id travels with the run instead: both run routes, `POST /v1/execute` and `POST /v1/start`, put it on the run's metadata, and the worker binds it from there while it runs the run's workflow and activities, so the lines written inside them carry the same id. A line the worker writes outside those bindings does not, such as the one the orchestration SDK logs after an activity has failed. One path does not carry it at all yet: when `POST /v1/validate` dispatches its dry run to a worker, the lines that worker writes for it have no `request_id`. The value is the one echoed in the response's `X-Request-ID` header and in the problem document's `request_id` member, so a caller reporting a failure hands you the key to its log lines.

The remaining keys are what the error handlers attach to an `event: "api_error"` record:

Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ classifiers = [
dependencies = [
# The extras deliberately leave out `cli`: the runner selects the `json` log sink and a Rich-free
# pretty-print mode, so nothing it does on a request renders a terminal.
"pipelex[mistralai,anthropic,google,google-genai,bedrock,fal]==0.67.0",
"pipelex[mistralai,anthropic,google,google-genai,bedrock,fal]==0.68.0",
"fastapi>=0.118.0",
"pyjwt>=2.10.1",
"uvicorn>=0.37.0",
Expand Down
62 changes: 57 additions & 5 deletions tests/unit/test_execute_dispatch.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@
JSON-safe output back into the full PipeOutput the `/execute` response wraps — exercising the real
serialize -> rehydrate round-trip (`serialize_completed_output` -> `hydrate_working_memory`), including
the `graph_spec` and `pipe_io_artifacts` `strict=False` re-validation branches. Also pins the policy-gated per-request override
(symmetric with `/start`) and the no-orchestrator case (`MissingOrchestratorError`). The boot slot is
(symmetric with `/start`), the no-orchestrator case (`MissingOrchestratorError`), and the request id reaching the job the
orchestrator is handed, which is the payload a Temporal worker deserializes and binds its log context from. The boot slot is
never used — every mode dispatches through the per-call registry. (Delivery is endpoint-set, never
requestable, so `/execute` has no fire-and-forget refusal — that axis is `/start`'s.)
"""
Expand Down Expand Up @@ -40,6 +41,7 @@

from api.api_config import ApiConfig
from api.exception_handlers import register_exception_handlers
from api.middleware import REQUEST_ID_HEADER, RequestIdMiddleware
from api.routes import router as api_router
from api.routes.pipelex.pipeline import ApiRunner
from tests.unit._constants import VALID_MTHDS
Expand All @@ -53,7 +55,8 @@ class _StubOrchestrator:
Returning via `serialize_completed_output` is the point — it produces the real JSON-safe
`PipelexPipeRunOutput` (the same shape that crosses the Temporal worker boundary), so the route
exercises the production serialize -> rehydrate round-trip instead of a hand-built payload. It
records each dispatch so a test can assert `/execute` drove the blocking `execute` arm. `start` (the
records each dispatch so a test can assert `/execute` drove the blocking `execute` arm, and what the
dispatched job's `RunMetadata` carries, since that is all a worker ever learns of the request. `start` (the
fire-and-forget arm) is present only to satisfy the protocol — `/execute` never calls it.
"""

Expand All @@ -72,7 +75,13 @@ def __init__(
self.supports_fire_and_forget = supports_fire_and_forget

async def execute(self, *, pipe_job: PipeJob, delivery_assignment: DeliveryAssignment | None) -> PipelexPipeRunOutput:
self.calls.append({"pipe_code": pipe_job.pipe.code, "delivery_assignment": delivery_assignment})
self.calls.append(
{
"pipe_code": pipe_job.pipe.code,
"delivery_assignment": delivery_assignment,
"request_id": pipe_job.job_metadata.run_metadata.request_id,
}
)
# A completed run always delivers a main stuff (pipelex invariant; enforced by
# `resolve_main_stuff_root_key` in both `serialize_completed_output` and `from_pipe_output`).
# This stub is an echo: promote the job's input stuff to the run's main stuff via the
Expand Down Expand Up @@ -129,11 +138,12 @@ def _echo_pipe_io_artifacts() -> PipeIOArtifacts:
)


def _build_client() -> TestClient:
def _build_client(*, with_request_id_middleware: bool = False) -> TestClient:
"""Wire the real routes; `with_request_id_middleware` wraps the app the way `api.main` does."""
app = FastAPI()
app.include_router(api_router, prefix="/v1")
register_exception_handlers(app)
return TestClient(app)
return TestClient(RequestIdMiddleware(app) if with_request_id_middleware else app)


def _register_stub(mocker: MockerFixture, *, mode: str, stub: _StubOrchestrator) -> None:
Expand Down Expand Up @@ -307,6 +317,48 @@ def test_forbidden_orchestration_mode_override_is_a_403(self) -> None:
assert response.headers["content-type"].startswith("application/problem+json")
assert response.json()["error_type"] == "OrchestrationModeOverrideForbidden"

def test_inbound_request_id_reaches_the_dispatched_job(self, mocker: MockerFixture) -> None:
"""On a `temporal` deployment, the inbound `X-Request-ID` rides the job the orchestrator is handed.

A worker never sees the request: it binds its log context from the deserialized job's
`RunMetadata.request_id`, and the middleware's in-process binding does not cross to it. So
the proof is the payload, not the runner call — a route that passed the id to a runner which
dropped it would still leave every worker line of the run without one.
"""
_force_config(mocker, mode="temporal", allow_override=False)
stub = _StubOrchestrator()
_register_stub(mocker, mode="temporal", stub=stub)
inbound_request_id = "01HNJZ4XR7K3Q9D8MWAQ7FY2E5"

client = _build_client(with_request_id_middleware=True)
response = client.post(
"/v1/execute",
json={"pipe_code": "echo", "mthds_contents": [VALID_MTHDS], "inputs": {"text": "hello"}},
headers={REQUEST_ID_HEADER: inbound_request_id},
)

assert response.status_code == 200, response.text
assert response.headers[REQUEST_ID_HEADER] == inbound_request_id
assert len(stub.calls) == 1
assert stub.calls[0]["request_id"] == inbound_request_id

def test_minted_request_id_reaches_the_dispatched_job(self, mocker: MockerFixture) -> None:
"""Without an inbound header, the job carries the id the middleware minted, which the response echoes."""
stub = _StubOrchestrator()
_register_stub(mocker, mode="direct", stub=stub)

client = _build_client(with_request_id_middleware=True)
response = client.post(
"/v1/execute",
json={"pipe_code": "echo", "mthds_contents": [VALID_MTHDS], "inputs": {"text": "hello"}},
)

assert response.status_code == 200, response.text
minted_request_id = response.headers[REQUEST_ID_HEADER]
assert minted_request_id
assert len(stub.calls) == 1
assert stub.calls[0]["request_id"] == minted_request_id

@pytest.mark.asyncio
async def test_missing_orchestrator_for_resolved_mode_raises(self, mocker: MockerFixture) -> None:
"""A resolved mode with no registered orchestrator fails loud with MissingOrchestratorError."""
Expand Down
20 changes: 18 additions & 2 deletions tests/unit/test_pipeline_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -334,8 +334,8 @@ def test_start_propagates_request_id_to_runner(self, mocker: MockerFixture):
# The middleware stores the inbound `X-Request-ID` on `request.state`; the route reads it
# back via `request_id_of(request)` and passes it as
# `request_id=` to `ApiRunner.start`, which forwards it to
# `pipeline_run_setup(...)` so it lands on `JobMetadata.request_id`.
# Without this hop the worker's `WorkflowLog` would carry `None`.
# `pipeline_run_setup(...)` so it lands on `RunMetadata.request_id`.
# Without this hop the worker's lines for the run would carry no `request_id`.
client, _, start_mock = _build_client(mocker, with_request_id_middleware=True)
inbound_request_id = "01HNJZ4XR7K3Q9D8MWAQ7FY2E5"
response = client.post(
Expand All @@ -348,6 +348,22 @@ def test_start_propagates_request_id_to_runner(self, mocker: MockerFixture):
start_mock.assert_awaited_once()
assert start_mock.await_args.kwargs["request_id"] == inbound_request_id

def test_execute_propagates_request_id_to_runner(self, mocker: MockerFixture):
# The `/execute` twin of the `/start` hop above: `ApiRunner.execute` forwards the id to the
# runtime's `execute`, which puts it on `RunMetadata.request_id`. On a Temporal deployment
# the whole run happens on a worker, which reads the id from that payload and nowhere else.
client, execute_mock, _ = _build_client(mocker, with_request_id_middleware=True)
inbound_request_id = "01HNJZ4XR7K3Q9D8MWAQ7FY2E5"
response = client.post(
"/v1/execute",
json={"pipe_code": "echo", "mthds_contents": [VALID_MTHDS], "inputs": {"text": "hello"}},
headers={REQUEST_ID_HEADER: inbound_request_id},
)
assert response.status_code == 200
assert response.headers[REQUEST_ID_HEADER] == inbound_request_id
execute_mock.assert_awaited_once()
assert execute_mock.await_args.kwargs["request_id"] == inbound_request_id

@pytest.mark.parametrize(
"bad_url",
[
Expand Down
Loading
Loading