From e16762036bb280a7d1510b68ff39d498b794cc59 Mon Sep 17 00:00:00 2001 From: Filippo Pacifici Date: Wed, 16 Sep 2026 09:28:10 -0700 Subject: [PATCH 1/3] Architecture docs --- AGENTS.md | 7 + README.md | 70 ++++- sentry_streams/AGENTS.md | 4 + sentry_streams/README.md | 11 + sentry_streams/docs/architecture/README.md | 187 +++++++++++++ sentry_streams/docs/architecture/contracts.md | 64 +++++ .../docs/architecture/guarantees.md | 69 +++++ .../docs/architecture/guides/adding-a-step.md | 253 ++++++++++++++++++ .../architecture/pipeline-dsl-and-runner.md | 136 ++++++++++ .../docs/architecture/rust-arroyo/README.md | 135 ++++++++++ .../docs/architecture/rust-arroyo/adapter.md | 82 ++++++ .../docs/architecture/rust-arroyo/messages.md | 184 +++++++++++++ .../docs/architecture/rust-arroyo/metrics.md | 61 +++++ .../rust-arroyo/python-operator.md | 139 ++++++++++ .../docs/reference/deployment-config.md | 184 +++++++++++++ .../deployment_config/README.md | 37 +-- 16 files changed, 1590 insertions(+), 33 deletions(-) create mode 100644 sentry_streams/docs/architecture/README.md create mode 100644 sentry_streams/docs/architecture/contracts.md create mode 100644 sentry_streams/docs/architecture/guarantees.md create mode 100644 sentry_streams/docs/architecture/guides/adding-a-step.md create mode 100644 sentry_streams/docs/architecture/pipeline-dsl-and-runner.md create mode 100644 sentry_streams/docs/architecture/rust-arroyo/README.md create mode 100644 sentry_streams/docs/architecture/rust-arroyo/adapter.md create mode 100644 sentry_streams/docs/architecture/rust-arroyo/messages.md create mode 100644 sentry_streams/docs/architecture/rust-arroyo/metrics.md create mode 100644 sentry_streams/docs/architecture/rust-arroyo/python-operator.md create mode 100644 sentry_streams/docs/reference/deployment-config.md diff --git a/AGENTS.md b/AGENTS.md index d54db8f9..87d489ad 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -16,6 +16,13 @@ The repository is organized into two main Python packages: - **`sentry_streams_k8s/`** - Kubernetes integration and deployment automation for Sentry Streams (pure Python package) - See [sentry_streams_k8s/AGENTS.md](./sentry_streams_k8s/AGENTS.md) for package-specific development instructions +## Architecture documents + +Overall structure of the streaming platform is in the [README](./README.md). +It contains links to the specific module. + +The architecture of the platform itself — subsystems, decisions, invariants, guarantees +and known debt — is in [sentry_streams/docs/architecture](./sentry_streams/docs/architecture/README.md). ### Other Directories diff --git a/README.md b/README.md index ea3556ba..2cb5b565 100644 --- a/README.md +++ b/README.md @@ -21,9 +21,6 @@ The main features are: - Support for stateful and stateless transformations. The state storage is provided by the platform rather than being part of the application. -- Distributed execution. The primitives used to build the application can - be distributed on multiple nodes by configuration. - - Hide the Kafka details from the application. Like commit policy and topic partitioning. @@ -34,7 +31,72 @@ The main features are: - Support for multiple runtimes. -[Streams Documentation](https://getsentry.github.io/streams/) +[Streams User Documentation](https://getsentry.github.io/streams/) + +## Architecture of a Streaming Application + +The three main components of a streaming application built with this library are: + +* The pipeline itself written via the Pipeline DSL. This represents the application logic. + +* The runtime. This is the binary that executes the pipeline primitives and processes data. + +* The Kubernetes infrastructure to deploy the pipeline (when on Kubernetes). + +```mermaid +flowchart TB + Kafka[Kafka] + + subgraph StreamingApp["Streaming Application"] + subgraph deployment["Kubernetes deployment"] + subgraph Pipeline["Pipeline"] + direction LR + source[source] --> transform[transform] --> sink[sink] + end + + Runtime[Runtime] + Runtime -->|executes| Pipeline + + Config[Pipeline ConfigMap] + end + + subgraph K8s["Kubernetes Infra"] + Macro["sentry-kube macro"] + Operator[Operator] + end + end + K8s -- manages --> deployment + K8s -- manages --> Config + Kafka --> source +``` + +Pipeline and runtime architecture docs are [here](./sentry_streams/docs/architecture/README.md) + +The pipeline DSL allows the user to define the streaming pipeline through a set +of dataflow primitives chained together. This is a Python DSL now. + +The application logic is separate from the infrastructure configuration which +is provided as a separate yaml file. The config file covers aspects like +observability, scale, parallelism, tuning, etc. This is meant to separate +the application logic from infra. + +The Runtime consumes the pipeline defined above and the config file, manages +the connectivity with Kafka and executes the processing in a supposedly optimized +way. + +The platform is provided as a Python package that contains also a native +portion as the runner. This package (`sentry_streams`) is imported by the application +code. The application defines the pipeline in its own code base and then uses +the runner provided by the library for the execution. + +This system provides the infrastructure needed to run the application in Kubernetes. +There are two options: + +* A [Sentry Kube](https://github.com/getsentry/sentry-infra-tools#sentry-kube) macro. + This can be used to generate Deployments and Configmap manifests to deploy manually + +* An experimental operator that manages a pipeline as a CRD. + ## to develop in this repo diff --git a/sentry_streams/AGENTS.md b/sentry_streams/AGENTS.md index ddbbeb78..9f0ad58d 100644 --- a/sentry_streams/AGENTS.md +++ b/sentry_streams/AGENTS.md @@ -24,6 +24,10 @@ the runner CLI script and the Arroyo adapter to run the application on single no - **`integration_tests/`** - Integration tests - **`docs/`** - Sphinx documentation +## Architecture documents + +Architecture is described [here](./docs/architecture/README.md) + ## Quick Start **Recommended:** Use the repository root Makefile commands documented in [../AGENTS.md](../AGENTS.md): diff --git a/sentry_streams/README.md b/sentry_streams/README.md index 388c4a2f..086329c8 100644 --- a/sentry_streams/README.md +++ b/sentry_streams/README.md @@ -6,3 +6,14 @@ is designed to handle real-time unbounded data streams. This is built primarily to allow the creation of Sentry ingestion pipelines though the api provided is fully independent from the Sentry product and can be used to build any streaming application. + +## Documentation + +- **[Architecture](./docs/architecture/README.md)** — the subsystems, the decisions behind + them, the rules that must hold, and what is unfinished. Start there whether you are building + on the platform or taking ownership of it. +- [Adding a step](./docs/architecture/guides/adding-a-step.md) — the procedure for extending + the DSL and the runtime. +- [Deployment configuration](./docs/reference/deployment-config.md) — the config file format. +- [User documentation](https://getsentry.github.io/streams/) — the published Sphinx docs. +- [AGENTS.md](./AGENTS.md) — development environment, tests, type checking. diff --git a/sentry_streams/docs/architecture/README.md b/sentry_streams/docs/architecture/README.md new file mode 100644 index 00000000..d73cc4f3 --- /dev/null +++ b/sentry_streams/docs/architecture/README.md @@ -0,0 +1,187 @@ +# Sentry Streams Architecture + +This directory describes the architecture of the `sentry_streams` package: the roles of +the main subsystems, the decisions behind them, the rules that must hold, and what we +know is unfinished. + +# The model + +This is the conceptual model the rest of the documentation assumes. Vocabulary is in +the [glossary](./README.md#glossary). + +## The three layers + +A streaming application built on this platform is made of three independent layers. +The whole design of the package follows from keeping them separate. + +**The pipeline DSL** is how a product team describes what their application does. It is +a pure, declarative description of a dataflow graph: a source, a series of +transformations, one or more sinks. It contains no reference to Kafka brokers, thread +counts, consumer groups, or the engine that will run it. An application is a Python +file that builds a `Pipeline` object and assigns it to a module level variable named +`pipeline`. + +**The adapter** turns that description into something executable on a specific runtime. +Each adapter implements a common interface, one method per primitive, and is free to +build whatever runtime-specific structure it wants as those methods are called. + +**The runner** is the glue and the entry point. It loads the application file and the +deployment configuration, instantiates the requested adapter, walks the pipeline graph +handing each step to the adapter, and tells the adapter to run. + +Alongside these there is **the deployment configuration**: a YAML file carrying +everything the pipeline description deliberately left out. It is described once, in +[Deployment configuration](../reference/deployment-config.md). + +```mermaid +flowchart TB + subgraph authored["Authored by the product team"] + app["application.py
(Pipeline DSL)"] + cfg["config.yaml
(deployment config)"] + end + + runner["Runner
sentry_streams/runner.py"] + + subgraph adapters["Adapters"] + rust["RustArroyoAdapter"] + py["ArroyoAdapter (Python)"] + other["...custom adapter"] + end + + runtime["Running consumer"] + + app -- "pipeline graph" --> runner + cfg -- "infrastructure config" --> runner + runner -- "translate_step() per step" --> adapters + adapters -- "run()" --> runtime +``` + +The key property of this split is that the same application file can be deployed +against different brokers, topics, parallelism settings and even different runtimes +without being modified. + +A fourth layer, the Kubernetes deployment automation, lives in the `sentry_streams_k8s` +package. The seam between the two is described in +[Deployment configuration](../reference/deployment-config.md#the-seam-with-sentry_streams_k8s). + +## The primitives + +The DSL is built from a small set of primitives, in a few families: + +- **Sources** produce messages into the pipeline. `StreamSource` reads from a Kafka + topic. A pipeline has exactly one root source. +- **Transforms** are 1:1 or 1:N operations: `Map`, `PredicateFilter`, `HeadersFilter`, + `FlatMap`. +- **Reduces** accumulate multiple messages into one over a window: `Batch` for plain + batching, `Aggregate` for windowed, optionally keyed, accumulator based aggregation. +- **Branching steps** split the stream: `Router` sends each message to exactly one + downstream branch, `Broadcast` sends a copy to every branch. Both take fully defined + sub-pipelines built with `branch()`. +- **Sinks** terminate a branch: `StreamSink` produces to Kafka, `GCSSink` writes objects + to GCS, `DevNullSink` discards. Every branch must end in a sink; the runner validates + this before anything is built. +- **Complex steps** are higher level primitives that are syntactic sugar over the simple + ones: `Parser`, `Serializer`, `BatchParser`, `ParquetSerializer`, `Reducer`. See + [Pipeline DSL and runner](./pipeline-dsl-and-runner.md#complex-steps). + +Refer to [pipeline.py](../../sentry_streams/pipeline/pipeline.py) to find documentaiton +on each of the implemented steps. + +A pipeline is built by chaining these together. Each step has a name, which is also the +key the deployment config uses to configure that step. + +```python +from sentry_kafka_schemas.schema_types.ingest_metrics_v1 import IngestMetric + +from sentry_streams.examples.transform_metrics import filter_events, transform_msg +from sentry_streams.pipeline.pipeline import ( + Map, + Parser, + PredicateFilter, + Serializer, + StreamSink, + streaming_source, +) + +pipeline = streaming_source(name="myinput", stream_name="ingest-metrics") + +( + pipeline.apply(Parser[IngestMetric]("parser")) + .apply(PredicateFilter("filter", function=filter_events)) + .apply(Map("transform", function=transform_msg)) + .apply(Serializer("serializer")) + .sink(StreamSink("mysink", stream_name="transformed-events")) +) +``` + +Notice what is absent: no broker address, no consumer group, no parallelism, no runtime. +`ingest-metrics` and `transformed-events` are logical stream names the deployment config +can override with the physical topics of a given environment. + +This snippet is the canonical example. Other documents refer to it rather than repeating +it. + +## Build, then run + +Building and running are two distinct phases, and the adapter interface is shaped around +that distinction. + +During the **build phase** the runner traverses the pipeline graph from the root source, +calling one adapter method per step through a `RuntimeTranslator`. The traversal follows +data flow rather than a list: each method receives the handle of the stream it attaches +to and returns the handle it produces, and branching steps return one handle per branch. +Nothing executes; the adapter accumulates. + +During the **run phase** the runner calls `run()` on the adapter, which blocks until the +consumer shuts down. + +That these are two separate functions is an invariant, not an implementation detail — +see [Contracts](./contracts.md#the-runner). + +## Multiple adapters + +The adapter interface is generic in the type of the stream handle, so adapters can +represent "a stream" however their runtime requires. The runner knows nothing about what +flows through it. + +The adapter used in production is the **Rust Arroyo adapter** (`rust_arroyo`), documented +in [The Rust Arroyo adapter](./rust-arroyo/README.md). + + +# Reading paths + +**Evolving the platform** + +2. [Adding a step](./guides/adding-a-step.md) — the procedure, with anchors and a skeleton. +3. [Contracts and invariants](./contracts.md) — the rules a new step must not break. +4. The [Rust runtime](./rust-arroyo/README.md) page for the area you are touching. + +**Architecture** + +2. [Guarantees and failure modes](./guarantees.md) — what the system promises. + + +# Glossary + +Several of these words are used in more than one sense elsewhere in the industry, and +two of them are used in more than one sense in this repository. + +| Term | Meaning here | +| --- | --- | +| **Step** | A node in the pipeline the application author writes: `Map`, `Batch`, `StreamSink`. | +| **Primitive** | A step type the adapter interface has a method for (`source`, `map`, `filter`, `reduce`, `sink`, `router`, `broadcast`, `flat_map`). | +| **Complex step** | A step that is syntactic sugar: it `convert()`s to simple steps before an adapter sees it. `Parser`, `Serializer`, `Reducer`. | +| **Adapter** | The component that turns a pipeline description into something a specific runtime can execute. | +| **Handle** | Whatever an adapter uses to represent "a stream" while the graph is being walked. Opaque to the runner. For the Rust adapter it is a `Route`. | +| **Route** | A source name plus an ordered list of waypoints. Identifies a branch of the pipeline. | +| **Waypoint** | One branch name appended to a route by a `Router` or `Broadcast`. | +| **Operator** (`RuntimeOperator`) | A Rust-side *descriptor* of something to build. One operator may cover several DSL steps. | +| **Strategy** | An Arroyo `ProcessingStrategy`. What an operator is built into at run time. | +| **Chain** / **fusion** | Consecutive `Map` steps accumulated by the adapter and composed into a single function. | +| **Delegate** | A Python object implementing `RustOperatorDelegate`, driven by the Rust `PythonAdapter` strategy. | +| **Segment** (config list) | An entry in `pipeline.segments`, selected with `--segment-id`. The deployment-level unit: one Kubernetes workload per entry. | +| **Segment** (`starts_segment`) | A boundary between chains *within* one process. Marks where fusion stops and a new parallelism setting applies. | +| **Watermark** | A periodic control message carrying accumulated offsets. Commits are driven by watermarks, not by data messages. | +| **Committable** | The set of `(topic, partition) → offset` entries a message or watermark is responsible for. | +| **`RawMessage`** | A payload the runtime understands: bytes, readable from Rust without the GIL. | +| **`PyAnyMessage`** | A payload only the application understands: an opaque reference to a Python object. | diff --git a/sentry_streams/docs/architecture/contracts.md b/sentry_streams/docs/architecture/contracts.md new file mode 100644 index 00000000..77d77015 --- /dev/null +++ b/sentry_streams/docs/architecture/contracts.md @@ -0,0 +1,64 @@ +# Contracts and invariants + +Rules that hold across the platform. Each one is load-bearing: something breaks, usually +subtly, if it stops holding. Use this as a review checklist. + +Behavioural guarantees to the outside world — delivery semantics, failure policy — are in +[Guarantees](./guarantees.md). + +## The runner + +| Invariant | Why | +| --- | --- | +| Loading the runtime and running it are separate functions | The Rust CLI embeds a Python interpreter, calls the loading function to obtain the built runtime, and calls `run()` itself from Rust. Merging them breaks that entry point | +| Exceptions from product code propagate | A Sentry SDK initialised by the application has to be able to capture them | +| Every branch terminates in a sink, validated before anything is built | A dangling branch is nearly always a bug; catching it after startup is far more expensive | + +## Configuration + +| Invariant | Why | +| --- | --- | +| `override_config()` then `validate()`, called by the adapter at translation time, and nowhere else | Before it, the step holds what the application wrote; after it, what this deployment wants. A step valid as written may be invalid once overridden, so validation must follow the override | +| No other component reads config keys off a step | One place to look when a value is wrong | +| Configuration is plain data, re-applied per process | It has to survive being pickled into a multiprocessing worker. A live backend object would not | +| Only non-semantic values are overridable | Functions, graph shape and types belong to the application; topics, brokers, consumer groups, batch sizes and parallelism belong to the deployment | + +## The adapter + +| Invariant | Why | +| --- | --- | +| The fused function of a closed chain is checked for picklability, even when multiprocessing is off | Otherwise a pipeline works in a single-process deployment and fails when someone enables a pool in production | +| Parallelism is configured only on a step that starts a segment | Parallelism belongs to the segment; a step in the middle of a chain is not its own execution unit | +| Nothing executes during the build phase | The graph traversal only accumulates; the runtime is assembled in `run()` | +| Operator descriptors are kept, not consumed | The strategy chain is rebuilt from them on every rebalance | + +## The message model + +| Invariant | Why | +| --- | --- | +| A strategy checks the route first and forwards untouched on mismatch, without taking the GIL | Every message physically traverses every strategy; this is what makes that affordable | +| A strategy forwards watermarks unless it buffers messages, in which case it holds them until the buffered work is out | A committed offset must be a processed offset | +| A `PyWatermark` is never submitted back into a Python operator | Directional conversion invariant of `WatermarkMessage` | +| A `PyWatermark` is never handed to the Kafka sink | Same | +| A step handles every payload variant: supports what it can, panics loudly on what it cannot | Silently forwarding an unsupported payload corrupts data downstream instead of failing at the bug | +| Payloads are not converted unless they are read | An opaque payload is a pointer move; a converted one is a copy plus a GIL acquisition | +| Messages are immutable; replacing a payload produces a new message | Enforced by Rust, a convention in Python | + +## Extending the DSL + +Adding a primitive is a **breaking change to every adapter**, by construction: the adapter +interface has one abstract method per primitive, and `RuntimeTranslator` dispatches on a +closed `StepType` enum. That is a deliberate trade. + +What it obliges you to do: + +1. Every adapter must implement the new method. An adapter that will not support the + primitive raises `NotImplementedError` — that is the sanctioned way to decline, and is + how `flat_map` is handled today in the Rust adapter. +2. `RuntimeTranslator.translate_step` must gain a branch, or `assert_never` fails the type + check. +3. If the primitive is not a map, its adapter methods must close the chain (see above). +4. Prefer a `ComplexStep` if the behaviour can be composed from existing primitives. It + costs no adapter changes and adapters may still override it natively. + +The procedure is in [Adding a step](./guides/adding-a-step.md). diff --git a/sentry_streams/docs/architecture/guarantees.md b/sentry_streams/docs/architecture/guarantees.md new file mode 100644 index 00000000..1abb92ca --- /dev/null +++ b/sentry_streams/docs/architecture/guarantees.md @@ -0,0 +1,69 @@ +# Guarantees and failure modes + +What the platform promises, and how it behaves when things go wrong. This is the page to +read before an incident, not during one. + +## Delivery semantics + +**At-least-once.** Offsets are committed only after the work they represent has been done: +the commit policy sits at the end of the strategy chain and commits the offsets a watermark +carries once that watermark has arrived from every branch. A crash between processing and +commit therefore replays. There is no deduplication and no transactional producer, so +duplicates are visible to sinks. + +Consequences worth stating plainly: + +- A Kafka sink may produce the same record more than once. This is not actually + working today. +- A GCS sink may write the same object more than once. +- Application steps must tolerate being run twice on the same input. + + +## Failure policy + +The runtime distinguishes three kinds of failure, deliberately: + +| Kind | Trigger | Behaviour | +| --- | --- | --- | +| **Invalid message** | A Python step raises `InvalidMessageError` (`pipeline/exception.py`); a delegate raises Arroyo's `InvalidMessage` | Offset and partition are read from the message and it is routed to the DLQ | +| **Backpressure** | A delegate raises `MessageRejected` | The message is handed back to the caller to be retried; pressure propagates upstream | +| **Bug** | Any other exception from a Python step; an unsupported payload variant reaching a step; a `PyAnyMessage` reaching the Kafka sink | Panic. Continuing would mean silently dropping or corrupting data | + +Two important qualifications: + +- **With no DLQ configured, an invalid message stops processing.** The consumer logs + `DLQ not configured, invalid messages will cause processing to stop` at startup. A DLQ is + configured per source via the `dlq` key. +- **An invalid message on an `AnyMessage` panics** rather than going to the DLQ, because + there is no offset to attribute it to. Messages produced downstream of the Python operator + are any-messages, which is why a failure there cannot be DLQ'd. + +DLQ limits are currently Arroyo's defaults — no rate limit and no cap on buffered messages. +That is a known gap, not a deliberate decision. + +## Backpressure + +Backpressure is Arroyo's: a strategy that cannot accept a message returns it, and the +processor stops polling until it can make progress. The Python operator participates +properly — `MessageRejected` from a delegate becomes Arroyo backpressure, and the drain of +`poll` results stops at the first rejection, preserving order. + +There is no load shedding and no queue with a bounded drop policy. A pipeline that cannot +keep up lags; it does not lose data. + +## Rebalance and shutdown + +- On **shutdown** (SIGINT/SIGTERM) the processor handle signals shutdown and the main loop + exits. `shutdown()` on the Python adapter is not implemented; shutdown is driven from Rust. +- On **rebalance**, Arroyo tears down and rebuilds the strategy chain. This is why the + consumer keeps `RuntimeOperator` descriptors rather than consuming them, and why delegates + are created by a factory: state too expensive to rebuild — notably the multiprocessing + pool — lives in the factory and survives. +- **Unverified:** behaviour when a multiprocessing worker dies mid-batch. + +## Health + +When `write_healthcheck` is enabled in the adapter config, the `HealthCheck` strategy +touches a file on every poll, which the Kubernetes liveness probe reads. The signal means +"the main loop is turning"; it does not mean the pipeline is making progress, and it will +keep succeeding while the consumer is backpressured. diff --git a/sentry_streams/docs/architecture/guides/adding-a-step.md b/sentry_streams/docs/architecture/guides/adding-a-step.md new file mode 100644 index 00000000..616a1101 --- /dev/null +++ b/sentry_streams/docs/architecture/guides/adding-a-step.md @@ -0,0 +1,253 @@ +# Adding a step + +End-to-end procedure for adding a new step to the pipeline DSL and implementing it in the +Rust runtime. Read [Messages → rules for step authors](../rust-arroyo/messages.md#rules-for-step-authors) +first; everything below assumes them. + +## 0. Decide what you are actually adding + +Cheapest first. Only go down the list if the option above genuinely cannot express it. + +```mermaid +flowchart TB + q1{"Can it be composed from
existing primitives?"} + q2{"Does it need to run
application (Python) code?"} + q3{"Does it buffer, or produce
on its own cadence?"} + q4{"Does a Python Arroyo strategy
already do it?"} + + complex["ComplexStep
convert() to simple steps.
No adapter changes."] + native["Native Rust strategy
+ RuntimeOperator variant."] + pycall["Rust strategy that calls
a Python callable
(like Filter)."] + delegate["RustOperatorDelegate
driven by PythonAdapter."] + wrap["ArroyoStrategyDelegate
wrapping it unmodified."] + + q1 -- yes --> complex + q1 -- no --> q2 + q2 -- no --> native + q2 -- yes --> q3 + q3 -- no --> pycall + q3 -- yes --> q4 + q4 -- no --> delegate + q4 -- yes --> wrap +``` + +A `ComplexStep` costs nothing outside the DSL and is the right answer more often than it +looks. A new **primitive** is the expensive case: it is a breaking change to every adapter +(see [Contracts](../contracts.md#extending-the-dsl)). Everything below covers that case. + +## Touchpoints + +| # | What | Where | +| --- | --- | --- | +| 1 | New `StepType` member | `sentry_streams/pipeline/pipeline.py:58` | +| 2 | Step dataclass | `sentry_streams/pipeline/pipeline.py` (next to `Map:427`, `Filter:444`) | +| 3 | DSL entry point, if not reachable via `apply()` | `Pipeline.apply:128`, `.sink:136`, `.broadcast:142`, `.route:154` | +| 4 | Abstract method on the adapter interface | `sentry_streams/adapters/stream_adapter.py:44` | +| 5 | Dispatch branch | `RuntimeTranslator.translate_step`, `stream_adapter.py:173` | +| 6 | Implementation in **every** adapter | `adapters/arroyo/rust_arroyo.py:288+`, `adapters/arroyo/adapter.py`, `dummy/dummy_adapter.py` | +| 7 | New `RuntimeOperator` variant | `src/operators.rs:33` | +| 8 | Arm in the operator dispatch | `src/operators.rs:125` (`build`) | +| 9 | The strategy itself | new `src/.rs`, registered in `src/lib.rs` | +| 10 | Tests | `tests/` (Python), `#[cfg(test)]` in your module (Rust) | + +You do **not** need to touch `sentry_streams/config.json`: `steps_config` allows additional +properties, so new per-step config keys are legal without a schema change. + +## 1. Choose the base class + +| Base | Use for | +| --- | --- | +| `Transform[TIn, TOut]` | A 1:1 or 1:N transformation | +| `Filter[TIn]` | Same type in and out, may drop | +| `Sink[TIn]` | Terminates a branch | +| `Source[TOut]` | Roots a pipeline | +| `WithInput[TIn]` | Branching or control steps that fit none of the above | +| `ComplexStep[TIn, TOut]` | Sugar over the above; implement `convert()` | + +## 2. Declare the step in the DSL + +```python +@dataclass +class Dedupe(Transform[TIn, TIn], Generic[TIn]): + """One-line docstring: this is the reference documentation for the step.""" + + window_size: int = 1000 + step_type: StepType = StepType.DEDUPE + + def override_config(self, loaded_config: Mapping[str, Any]) -> None: + if "window_size" in loaded_config: + self.window_size = loaded_config["window_size"] + + def validate(self) -> None: + if self.window_size <= 0: + raise InvalidPipelineError(f"{self.name}: window_size must be positive") +``` + +Notes that save time: + +- It is a `@dataclass`, so fields with defaults must follow fields without them — including + the inherited ones. `step_type` always carries a default. +- Generics are load-bearing: `Pipeline[TOut]` is re-typed by `apply()`, so declaring the wrong + parameters shows up as a chaining type error in application code, not here. Run + `make typecheck`, which includes `tests/test_mypy_integration.py`. +- `override_config` reads only keys it understands; `validate()` runs after it, never before. +- Put the per-primitive documentation in the docstring. The architecture docs deliberately do + not duplicate it. + +## 3. Wire the translator + +Add the `StepType` member, then a branch in `RuntimeTranslator.translate_step`. The dispatch +ends in `assert_never(step_type)`, so mypy fails until you do — that is the intended +enforcement, not an obstacle to work around. + +## 4. Implement it in every adapter + +Add the abstract method to `StreamAdapter`, then implement it everywhere. An adapter that will +not support the step declines explicitly: + +```python +def dedupe(self, step: Dedupe[Any], stream: Route) -> Route: + logger.info(f"Adding dedupe: {step.name} to pipeline") + raise NotImplementedError +``` + +In `RustArroyoAdapter`, follow the five-step pattern in +[the adapter in depth](../rust-arroyo/adapter.md#what-each-method-does). The step that is easy +to forget is **closing the open transformation chain**, which every non-map method must do +before adding its own operator: + +```python +def dedupe(self, step: Dedupe[Any], stream: Route) -> Route: + step_config = self.__get_step_config(step.name) + step.override_config(step_config) + step.validate() + self.__close_chain(stream) # required — see Contracts + self.get_consumer(stream.source).add_step( + RuntimeOperator.Dedupe(stream, step.name, step.window_size) + ) + return stream +``` + +## 5. Implement the Rust strategy + +Add the `RuntimeOperator` variant (`src/operators.rs:33`) and its arm in `build` +(`src/operators.rs:125`), then write the strategy. This skeleton applies all six step-author +rules; `src/filter_step.rs` is the closest real example. + +```rust +use crate::messages::RoutedValuePayload; +use crate::pipeline_stats::get_stats; +use crate::routes::{Route, RoutedValue}; +use crate::utils::traced_with_gil; +use sentry_arroyo::processing::strategies::{ + CommitRequest, ProcessingStrategy, StrategyError, SubmitError, +}; +use sentry_arroyo::types::{InnerMessage, Message}; +use std::time::{Duration, Instant}; + +pub struct MyStep { + pub next_step: Box>, + pub route: Route, + pub step_name: String, +} + +impl ProcessingStrategy for MyStep { + fn poll(&mut self) -> Result, StrategyError> { + // If this step buffers: emit anything ready here, and release any watermark + // whose offsets the emitted output now covers (rule 2), before polling on. + self.next_step.poll() + } + + fn submit(&mut self, message: Message) -> Result<(), SubmitError> { + // Rule 1: not our branch — forward untouched, no GIL. + // Rule 2: forward watermarks too, UNLESS this step buffers messages, in which + // case hold them until the buffered work is out. + if self.route != message.payload().route || message.payload().payload.is_watermark_msg() { + return self.next_step.submit(message); + } + + // Rule 3: handle every payload variant you support; fail loudly on the rest. + let RoutedValuePayload::PyStreamingMessage(ref streaming_msg) = message.payload().payload + else { + unreachable!("watermarks are forwarded above") + }; + + let stats = get_stats(); + stats.step_exec(&self.step_name); + let start = Instant::now(); + + // Rules 4 and 5: only take the GIL if you must read a Python payload, + // and take it once around the largest reasonable unit of work. + // let outcome = traced_with_gil!(|py| { ... }); + + stats.step_timing(&self.step_name, start.elapsed().as_secs_f64()); + + // To send a message to the DLQ instead of forwarding it: + // + // match &message.inner_message { + // InnerMessage::BrokerMessage(broker_message) => { + // stats.step_error(&self.step_name); + // return Err(SubmitError::InvalidMessage(broker_message.into())); + // } + // // No offset to attribute the failure to — see Guarantees. + // InnerMessage::AnyMessage(..) => panic!("cannot DLQ an AnyMessage in {}", self.step_name), + // } + + self.next_step.submit(message) + } + + fn terminate(&mut self) { + self.next_step.terminate() + } + + fn join(&mut self, timeout: Option) -> Result, StrategyError> { + // If this step buffers: flush in-flight work and release held watermarks here. + self.next_step.join(timeout) + } +} +``` + +If the step needs to run a Python *step* rather than a Python function — it buffers, or has its +own cadence — do not write a strategy. Implement a `RustOperatorDelegate` instead and let +`RuntimeOperator::PythonAdapter` drive it; see +[The Python operator](../rust-arroyo/python-operator.md). + +## 6. Instrument it + +Use `get_stats()` (`src/pipeline_stats.rs`): `step_exec` on entry, `step_timing` with the +elapsed seconds, `step_error` when the step rejects a message. Python-side steps are wrapped +automatically by the adapter. Keeping both sides identical is the point — see +[Metrics](../rust-arroyo/metrics.md#deliberate-symmetry-across-the-boundary). + +## 7. Configure it + +Add the keys your step reads to `override_config`, and document them in +[Deployment configuration](../../reference/deployment-config.md#per-step-configuration). No +JSON-schema change is needed. If the step is not a map, remember it will close the enclosing +transformation chain, which can split a segment a user thought was contiguous. + +## 8. Test it + +| Level | Where | What it gives you | +| --- | --- | --- | +| Rust unit | `#[cfg(test)]` in your module | `fake_strategy::FakeStrategy` collects what you forward; `assert_messages_match` compares; `testutils::{build_routed_value, make_lambda, import_py_dep}` build inputs and Python callables | +| Python unit | `tests/pipeline/`, `tests/adapters/` | DSL construction, config overrides, translation | +| Type | `make typecheck` | Chaining types; `assert_never` exhaustiveness | +| Integration | `integration_tests/` | The step inside a running consumer | + +Write at least: route mismatch forwards untouched; a watermark behaves correctly (forwarded, +or held and released); each supported payload variant; and the unsupported variant fails the +way you intended. + +## Checklist + +- [ ] `StepType` member added, translator branch added, `assert_never` satisfied +- [ ] Step dataclass with docstring, `override_config`, `validate` +- [ ] Every adapter implements the method, or raises `NotImplementedError` +- [ ] Non-map adapter methods close the transformation chain first +- [ ] `RuntimeOperator` variant + `build` arm +- [ ] Route checked first, without the GIL +- [ ] Watermarks forwarded, or held and released on flush +- [ ] Every payload variant handled or loudly rejected +- [ ] `step_exec` / `step_timing` / `step_error` emitted +- [ ] Rust and Python tests, `make typecheck`, `make tests-streams`, `make tests-rust-streams` diff --git a/sentry_streams/docs/architecture/pipeline-dsl-and-runner.md b/sentry_streams/docs/architecture/pipeline-dsl-and-runner.md new file mode 100644 index 00000000..0c041242 --- /dev/null +++ b/sentry_streams/docs/architecture/pipeline-dsl-and-runner.md @@ -0,0 +1,136 @@ +# Pipeline DSL and runner + +The layer product teams write against, and the component that turns what they wrote into +a running process. The conceptual model is in [The model](./README.md#the-model); the per-primitive +reference is in the docstrings of `sentry_streams/pipeline/pipeline.py`. + +## Principles + +**The pipeline is a description, not a program.** Building a `Pipeline` object executes +no streaming logic and allocates no runtime resource. It produces a graph: a map of named +steps and the edges between them. + +**Steps are runtime agnostic.** A primitive describes an intent and never how the intent +is fulfilled. The same `Map` may end up as a native Rust `RunTask`, as a link in a fused +chain, or as a function executed in a pool of worker processes. Product code does not +change. + +**Names are the contract with the deployment config.** Every step has a name, unique +within the pipeline, and that name is the key the config addresses it by. This is the +only coupling between the application file and its configuration. + +**Steps are typed.** `Pipeline[TOut]` carries the type of the messages flowing out of the +last step added, and each primitive is generic in its input and output types. Chaining a +step whose input type does not match is a type error — which catches, for instance, +sinking parsed objects into a Kafka sink without a serializer in between. What this +requires of a new step class is in +[Adding a step](./guides/adding-a-step.md#2-declare-the-step-in-the-dsl). + +## Chaining + +A pipeline starts from a source and grows by chaining — see the canonical example in +[The model](./README.md#the-primitives). `apply()` registers a step and an edge from the +previous one, then returns the same pipeline object re-typed to the new output type. +`sink()` does the same and closes the pipeline: nothing can be appended after it. + +Branching works by building sub-pipelines separately and handing them to a branching +step: + +```python +branch_a = branch("a").apply(Map("map_a", function=f)).sink(StreamSink("sink_a", stream_name="a")) +branch_b = branch("b").apply(Map("map_b", function=g)).sink(StreamSink("sink_b", stream_name="b")) + +pipeline.route( + "router", + routing_function=pick_branch, + routing_table={BranchKey.A: branch_a, BranchKey.B: branch_b}, +) +``` + +A sub-pipeline built with `branch()` has a `Branch` step as its root rather than a source. +When handed to `route()` or `broadcast()`, its steps and edges are merged into the parent +graph with the branching step as the parent of each branch root. Because the branches must +be complete before they can be merged, a router or broadcast also closes the pipeline it +is added to. + +The result, in all cases, is a single `Pipeline` object holding a flat map of steps plus +the incoming and outgoing edges. Downstream, nothing needs to know in which order the +chaining calls were made. + +## Configuration overrides + +Every step declares defaults in code and can have them overridden by the deployment +config, through two hooks on the base `Step` class: + +- `override_config(loaded_config)` — the step picks the keys it understands out of the + config mapping for its own name and mutates itself accordingly. +- `validate()` — checks the step is coherent. Called explicitly *after* the override, + because a step that was valid as written may not be valid once overridden. + +Both are invoked by the adapter at translation time, immediately before the step becomes +a runtime primitive. That placement is an invariant: see +[Contracts](./contracts.md#configuration). + +What is overridable and what is not follows one line: the values that change the *meaning* +of the program — the functions, the graph shape, the types — belong to the application. +Everything about *where and how much* — topics, brokers, consumer groups, batch sizes, +parallelism — is overridable. Which keys each step reads is documented in +[Deployment configuration](../reference/deployment-config.md#per-step-configuration). + +## Complex steps + +Some primitives are not runtime operations at all, but recurring combinations of them. +`Parser` decodes bytes and validates them against the `sentry-kafka-schemas` schema for a +message type; `Serializer` does the reverse; `ParquetSerializer` turns a batch into a +Parquet buffer; `Reducer` is a friendlier spelling of `Aggregate`. + +These are `ComplexStep`s. Each implements `convert()`, returning a plain simple step — +usually a `Map` with a suitable function partial-applied. The translator calls `convert()` +transparently, so by the time the adapter sees the pipeline, a `Parser` is just a `Map`. + +The indirection does two things. It keeps the set of operations an adapter must implement +small: an adapter supports maps and filters and reduces, not parsers and serializers. And +it leaves room for an adapter to do better than the generic conversion, through +`complex_step_override()`: an adapter with a native, faster implementation of a complex +step declares it and receives the original step instead of its conversion. Adapters that +do not care return an empty mapping. + +## From pipeline to consumer + +The runner is deliberately thin: the intelligence about a runtime lives in its adapter, +and the intelligence about the application lives in the pipeline. + +```mermaid +sequenceDiagram + participant CLI as runner CLI + participant Cfg as load_config + participant Sub as subprocess + participant Ad as adapter + participant Tr as RuntimeTranslator + + CLI->>Cfg: read YAML, resolve ${envvar:...}, validate schema + CLI->>Sub: exec application file + Sub-->>CLI: Pipeline object + CLI->>CLI: validate_all_branches_have_sinks + CLI->>CLI: configure_metrics + CLI->>Ad: load_adapter(name, config, metrics, segment_id) + CLI->>Tr: iterate_edges(pipeline, translator) + loop for each step, following edges + Tr->>Ad: source / map / filter / reduce / sink / router / broadcast + Ad-->>Tr: stream handle(s), one per output branch + end + CLI->>Ad: run() +``` + +| Component | Anchor | Note worth knowing | +| --- | --- | --- | +| `load_config` | `pipeline/config.py` | Resolves `${envvar:NAME}`, validates against `sentry_streams/config.json` | +| `_load_pipeline` | `runner.py` | Runs the application file **out of process**, so module-scope imports cannot pollute the runner. Product exceptions propagate, so a Sentry SDK the application initialised can capture them | +| `validate_all_branches_have_sinks` | `pipeline/validation.py` | A dangling branch is nearly always a bug, and far cheaper to catch here than after startup | +| `load_adapter` | `adapters/loader.py` | Resolves the adapter name to a class; with a segment id, narrows the config first so an adapter only ever sees its own portion | +| `RuntimeTranslator` | `adapters/stream_adapter.py:162` | The only place mapping step type to adapter method. Ends in `assert_never`, so a new `StepType` is a type error until handled | +| `iterate_edges` | `pipeline/pipeline.py` | The data-flow traversal. Branching steps return several handles, each pushed back into the working set | +| `adapter.run()` | per adapter | Blocks until shutdown | + +The handle is opaque to the runner — a generic parameter of the adapter. For the Rust +Arroyo adapter it is a `Route`; for another runtime it could be a native stream object. diff --git a/sentry_streams/docs/architecture/rust-arroyo/README.md b/sentry_streams/docs/architecture/rust-arroyo/README.md new file mode 100644 index 00000000..e203ea41 --- /dev/null +++ b/sentry_streams/docs/architecture/rust-arroyo/README.md @@ -0,0 +1,135 @@ +# The Rust Arroyo adapter + +The `rust_arroyo` adapter runs a pipeline as a single Arroyo consumer implemented in +Rust. It is the adapter used in production. + +- [Messages](./messages.md) — the message model and the rules it imposes on steps. +- [The Python operator](./python-operator.md) — embedding Python strategies in the chain. +- [The adapter in depth](./adapter.md) — chaining, fusion, segments. +- [Metrics](./metrics.md) — how both sides report. + +## Execution model + +The consumer is written in Rust on top of `rust-arroyo` and exposed to Python as an +extension module built with PyO3 and Maturin. The process is a Python process: it starts +as Python, imports the extension module, and hands control to Rust. + +This gives an inversion of control that explains most of the rest of this directory: + +- **Python builds.** The adapter constructs a description of the pipeline and hands it to + Rust step by step. +- **Rust runs.** Once `run()` is called, the Rust `StreamProcessor` owns the main loop: it + polls Kafka, submits messages through the strategy chain, and commits. +- **Rust calls back into Python.** The application logic — map functions, filter + predicates, routing functions — lives in Python memory. Every time a message reaches one + of those steps, Rust acquires the GIL and calls into Python. + +```mermaid +flowchart LR + subgraph py["Python"] + adapter["RustArroyoAdapter"] + appfn["Application functions"] + delegates["Python delegates
(multiprocess, reduce)"] + end + + subgraph rs["Rust"] + consumer["ArroyoConsumer"] + chain["Arroyo strategy chain"] + end + + adapter -- "add_step(RuntimeOperator)" --> consumer + adapter -- "run()" --> consumer + consumer -- "builds" --> chain + chain -- "call via PyO3 (GIL)" --> appfn + chain -- "submit / poll" --> delegates +``` + +## Building the consumer + +`ArroyoConsumer` is a `pyclass`. The adapter creates one per source, then calls +`add_step()` once per operator, in the order a message would traverse them. + +The values passed to `add_step()` are `RuntimeOperator`s (`src/operators.rs`): a Rust enum, +also exported to Python, whose variants are the operations the runtime knows how to +execute. They are **descriptors, not strategies** — they say what to build and carry the +parameters, including a reference to a Python callable where relevant. + +`RuntimeOperator`s do not map one-to-one onto DSL steps; the adapter fuses consecutive maps +into one operator. See [the adapter in depth](./adapter.md). + +The whole pipeline must be handed over before anything is constructed, because Arroyo +strategies are built back to front — each strategy owns the next. That is why the consumer +accumulates a list of operators and only assembles the chain in `run()`, and why it must +keep the descriptors rather than consuming them: the chain is rebuilt from them on every +rebalance, through a strategy factory. + +## What runs where + +The rule governing this split is **any primitive that can be executed entirely in Rust is +implemented in Rust**, because crossing the language boundary costs a GIL acquisition, +serialises the step against every other Python interaction in the process, and may cost a +copy. + +| Concern | Rust | Python | +| --- | --- | --- | +| Kafka consumption, offset management, rebalancing | ✅ | | +| Strategy chain, backpressure, DLQ routing | ✅ | | +| Watermark emission and the commit policy | ✅ | | +| Route matching and branch dispatch | ✅ | | +| Batch step, header filter, broadcast | ✅ | | +| Kafka producer sink, GCS sink, DevNull sink | ✅ | | +| Map and filter functions written by applications | invoked from | executed in | +| Routing functions | invoked from | executed in | +| Parsing, serialisation, schema validation | | ✅ (as maps) | +| Windowed reduce / aggregate | wrapper | strategy | +| Multiprocessing | wrapper | strategy | +| Building the pipeline, reading the config | | ✅ | +| Metrics emission | ✅ | ✅ | + +The two wrapper rows are pre-existing Python Arroyo strategies that were not +reimplemented; they are embedded through [the Python operator](./python-operator.md). + +There is a middle ground worth knowing about: a function written in Rust and exposed to the +pipeline as if it were a Python callable (the `rust_transforms` example). It is still +invoked through PyO3, but releases the GIL internally for the duration of the work. + +## The strategy chain + +When `run()` is called, the consumer builds an Arroyo strategy chain and hands it to a +`StreamProcessor` bound to the source topic. The chain is assembled from the end backwards; +a message traverses it in this order: + +```mermaid +flowchart TB + kafka["Kafka consumer"] + hc["HealthCheck
(optional)"] + conv["RunTask: to_routed_value
KafkaPayload → RoutedValue"] + wm["WatermarkEmitter"] + ops["Operators, in pipeline order
(one strategy per RuntimeOperator)"] + commit["WatermarkCommitOffsets"] + + kafka --> hc --> conv --> wm --> ops --> commit +``` + +| Stage | Role | +| --- | --- | +| **HealthCheck** | Touches a file on every poll for the Kubernetes liveness probe. Added only when enabled in the adapter config | +| **Conversion** | Turns the `KafkaPayload` into a `RoutedValue`: a payload plus the `Route` saying which branch it belongs to. The payload at this point is a `RawMessage` with the broker bytes, headers, timestamp and the logical schema name | +| **WatermarkEmitter** | Records the offsets of everything passing through and periodically injects a watermark carrying them | +| **Operators** | The pipeline proper. Each checks the message's route against its own and passes it straight through on a mismatch | +| **WatermarkCommitOffsets** | Commits the offsets a watermark carries, once it has seen that watermark arrive from every branch | + +A dead letter queue policy can be attached to the processor, in which case a message +rejected as invalid by a step is produced to the DLQ topic instead of stopping the +consumer. The full failure policy is in [Guarantees](../guarantees.md). + +The last two stages exist because Arroyo has no notion of branches. That is the subject of +[Messages](./messages.md). + +## Entry points + +Two executables converge here: the **Python CLI** (`sentry_streams.runner:main`), which +builds the runtime and calls `run()` from Python, and the **Rust CLI**, which embeds a +Python interpreter, calls the loading function to get the built runtime back as a +`Py`, and calls `run()` from Rust. Keeping loading and running separate is therefore +an [invariant](../contracts.md#the-runner). diff --git a/sentry_streams/docs/architecture/rust-arroyo/adapter.md b/sentry_streams/docs/architecture/rust-arroyo/adapter.md new file mode 100644 index 00000000..24fa22e2 --- /dev/null +++ b/sentry_streams/docs/architecture/rust-arroyo/adapter.md @@ -0,0 +1,82 @@ +# The adapter in depth + +`RustArroyoAdapter` (`sentry_streams/adapters/arroyo/rust_arroyo.py`) is the Python half of +the Rust runtime: the component the runner drives while walking the pipeline graph, which +ends up holding a fully described Rust consumer. + +## The stream handle is a route + +The adapter interface is generic in the type representing "a stream". This adapter uses +`Route` — the source name plus the waypoints identifying a branch. + +That choice says a lot about the runtime. There is no stream object to hand around: the +consumer is a single linear chain of strategies, and the only thing distinguishing one +logical stream from another is which branch it belongs to. So when the runner asks the +adapter to attach a map to a stream, what it passes is "which branch", and what it gets +back is the branch the output belongs to — the same one for everything except router and +broadcast, which return one route per branch. + +The adapter keeps a map of consumers keyed by source name. The source step creates the +consumer; every subsequent step looks up the consumer for its route's source and appends to +it. More than one consumer per process is not supported today, but the structure +anticipates it. + +## What each method does + +The per-step methods (`rust_arroyo.py:288` onwards) all follow one pattern: + +| # | Action | Note | +| --- | --- | --- | +| 1 | Look up this step's slice of `steps_config`, by step name | | +| 2 | `override_config()`, then `validate()` | The single point where deployment config is applied — an [invariant](../contracts.md#configuration) | +| 3 | Close any open transformation chain on this route, unless the step is a map | **Required.** Skipping it emits operators out of order — see [Contracts](../contracts.md#the-adapter) | +| 4 | Translate into one or more `RuntimeOperator`s and `add_step()` them | | +| 5 | Return the route, or the per-branch routes | | + +Four methods do more than the pattern: + +- **`source`** creates the `ArroyoConsumer` rather than adding to one: Kafka consumer + config, DLQ config if the source declares a DLQ stream, the Rust-side metrics config, the + healthcheck flag and the Sentry DSN. It captures the *logical* stream name before applying + overrides and passes that as the schema name, so a topic override does not change which + schema messages are validated against. +- **`filter`** distinguishes a `HeadersFilter`, which becomes a fully native operator with + no Python call, from a `PredicateFilter`, which carries a Python callable. +- **`reduce`** distinguishes `Batch`, which has a native Rust implementation, from + everything else, which is wrapped into a Python delegate and executed by the + [Python operator](./python-operator.md). +- **`router`** and **`broadcast`** add their operator and return one route per branch, each + with the branch name appended as a waypoint. That is what makes the runner's traversal + continue down every branch. + +Application functions are wrapped before registration so that entering the step, its +duration and any exception are recorded — see [Metrics](./metrics.md#deliberate-symmetry-across-the-boundary). + +`run()` asserts that exactly one consumer was built and calls `run()` on it, blocking in +Rust until shutdown. + +## Chaining and `starts_segment` + +The most interesting thing the adapter does is *not* translate steps one-to-one. The +mechanics are: + +| Rule | Detail | +| --- | --- | +| **What accumulates** | Consecutive `Map` steps, in a structure keyed by route (`steps_chain.py`), so each branch accumulates independently | +| **When a chain opens** | A map arrives and no chain is open on that route, or a map is marked `starts_segment: True` | +| **When a chain closes** | Any non-map step — filter, reduce, sink, router, broadcast — closes it *before* the new step is added, which keeps operator order faithful to pipeline order | +| **What a closed chain becomes** | With no parallelism, a native `RuntimeOperator::Map` holding the fused function. With `multi_process`, a `RuntimeOperator::PythonAdapter` wrapping Arroyo's multiprocessing strategy, whose pool executes the fused function | +| **Picklability check** | The fused function is checked for picklability when the chain is finalised, *even without multiprocessing*, so a pipeline cannot work in single-process and fail when someone enables a pool in production. Build-time error, naming the offending steps | + +Two consequences surprise people: + +- **Parallelism can only be configured on a step that starts a segment.** Configuring it + mid-chain is rejected: parallelism belongs to the segment, not the step. +- **A non-map step in the middle of what looks like one logical segment splits it.** + Inserting a filter between two maps produces two chains, two operators, and potentially + two pools. Segmentation is a consequence of the pipeline's shape as much as of the + configuration. + +Both senses of the word "segment" are defined in the +[glossary](../README.md#glossary); how they are expressed in the config file is in +[Deployment configuration](../../reference/deployment-config.md#segments). diff --git a/sentry_streams/docs/architecture/rust-arroyo/messages.md b/sentry_streams/docs/architecture/rust-arroyo/messages.md new file mode 100644 index 00000000..46c3d93f --- /dev/null +++ b/sentry_streams/docs/architecture/rust-arroyo/messages.md @@ -0,0 +1,184 @@ +# Messages + +The message model decides what a step can and cannot do, where the GIL has to be taken, +which conversions are possible, and what a step author has to handle. Almost every +limitation of this runtime traces back to something on this page. + +It is written rules-first: if you are implementing a step, the first section is the one +you must not skip, and the skeleton that applies it is in +[Adding a step](../guides/adding-a-step.md#5-implement-the-rust-strategy). + +## Rules for step authors + +1. **Check the route before anything else**, and forward on mismatch. No GIL. +2. **Handle watermarks explicitly.** Forward them unless the step buffers messages, in + which case hold them until the buffered work is out. +3. **Handle every payload variant.** Support what the step can support; panic with a clear + message on what it cannot, rather than silently passing it through. +4. **Do not convert payloads you do not need to read.** An opaque payload is a pointer + move; a converted one is a copy plus a GIL acquisition. +5. **Take the GIL once, in the largest reasonable unit of work.** +6. **Assume messages are immutable.** Replacing a payload produces a new message. This is + enforceable in Rust and not in Python — a Python step can mutate the object behind a + `PyAnyMessage` — so on the Python side it is a convention. + +The rest of this page is why. + +## Two classes of payload + +A message flowing through the Rust consumer has metadata — headers, a timestamp, an +optional schema name — and a payload. The metadata is the same for everyone; the payload +is what differs, in two fundamentally different ways. + +**Rust-native payloads** live in Rust memory and can be read without the GIL. Today the +only one is `RawMessage`, whose payload is a `Vec`: the bytes read from Kafka, or the +bytes a step produced to be written out. Rust can inspect, slice, hash or write it freely. + +**Python payloads** live in Python memory. `PyAnyMessage` holds a `Py` — a smart +pointer to an arbitrary Python object. From the Rust side it is *opaque*: it can be moved, +cloned (with the GIL) and handed back to Python, but not inspected. This is the +representation of everything between a parser and a serializer. + +The distinction states who owns the data: + +> A `RawMessage` is data the runtime understands. A `PyAnyMessage` is data the +> *application* understands, which the runtime only carries between the steps that do. + +Both are `pyclass`es, so both can be handed to Python code — the difference is what it +costs. Handing over a `PyAnyMessage` passes a pointer; handing over a `RawMessage` builds a +Python `bytes` object, which copies. + +## The enums + +``` +RoutedValue +├── route: Route which branch this message belongs to +└── payload: RoutedValuePayload + ├── PyStreamingMessage a data message + │ ├── PyAnyMessage payload is a Python object + │ └── RawMessage payload is a byte array + └── WatermarkMessage a control message + ├── Watermark native, Rust memory + └── PyWatermark crossed into Python memory +``` + +Every strategy in the chain is a `ProcessingStrategy`, so every strategy +receives every variant, and the compiler makes matching on them unavoidable. That is what +rules 1–3 are about. + +The "fail loudly" half of rule 3 is deliberate. The Kafka sink panics if it is given a +`PyAnyMessage`, because it has no way to turn an arbitrary Python object into bytes — that +is the serializer's job, and its absence from the pipeline is a bug in the application, +not a runtime condition to recover from. The general policy is in +[Guarantees](../guarantees.md#failure-policy). + +## Conversion is limited, on purpose + +The only conversion the runtime performs is **between byte arrays**: a `RawMessage` can be +exposed to Python as `bytes`, and bytes coming back from Python become a `RawMessage` +again. When a Python transform returns `bytes` the result is wrapped as a `RawMessage`; +when it returns anything else it is wrapped as a `PyAnyMessage`. + +There is no conversion from `PyAnyMessage` to `RawMessage`, because there is no general way +to serialise an arbitrary Python object. + +Two rules applications actually run into follow from that: + +- A pipeline that writes to Kafka must serialise before the sink. Without a serializer the + sink receives an opaque object and panics. +- A Rust-native step that needs to read the payload can only be applied where the payload + is a `RawMessage` — that is, before any step that produced Python objects. + +## Conversion costs the GIL + +Anything that touches a Python payload needs the GIL, including operations that look +innocuous. Cloning a `RawMessage` is a memcpy; cloning a `PyAnyMessage` is a `clone_ref`, +which increments a Python reference count and therefore requires attaching to the +interpreter — so a broadcast fanning one message out to four branches takes the GIL to do +it. Reading the timestamp is the same: the field lives inside a `pyclass` instance. + +Hence rules 4 and 5, and one further property: **GIL acquisition is instrumented.** All +Rust code acquires the GIL through the `traced_with_gil!` wrapper (`src/utils.rs`), which +logs a warning when acquisition takes longer than a threshold, because a slow acquisition +means another thread is holding it and the pipeline is serialising on Python. + +## Routes + +Arroyo pipelines are linear: a strategy has exactly one next strategy, and every message +visits every step in order. The DSL, however, has `Router` and `Broadcast`. Routes are how +the second is implemented on top of the first. + +Every message carries a `Route`: the name of the source it came from and an ordered list of +`waypoints`, the branches it has been sent down. Every strategy is built with the route it +belongs to, and forwards untouched anything whose route differs. The chain is therefore +still physically linear and contains the steps of every branch one after another; each +message "executes" only the steps whose route it matches. + +- A `Router` appends the one waypoint its routing function selected. +- A `Broadcast` emits one copy per downstream branch, each with a different waypoint. + +The cost is that a message physically traverses every strategy in the consumer, including +branches it will never enter. The route check is cheap and GIL-free, which is what makes +that acceptable. + +## Watermarks + +Branching breaks committing, and watermarks are the fix. + +Arroyo's normal commit policy commits the offsets of a message once it reaches the end of +the chain. With a broadcast there are several copies of that message in flight, and the +message is only really processed when *all* of them have finished; committing when the +first arrives would commit work that has not happened and lose it on a restart. + +A **watermark** is a control message injected by the `WatermarkEmitter` at the head of the +chain. The emitter records the offsets of every message passing through it and every few +seconds (10 by default) emits a watermark carrying the accumulated committable, then +clears it. + +Watermarks travel down the pipeline like data messages and go through the *same* branching +logic: a broadcast duplicates a watermark exactly as it duplicates a message, and a router +forwards one to all of its branches, because it cannot know which branches received +messages since the last one. + +At the end of the chain, the commit policy does not commit on messages at all. It counts +watermarks: for a given watermark it accumulates copies, and only when it has seen one from +every branch does it turn the carried offsets into a commit request. + +``` + Source + | +┎-Router-┓ +| | +1 ┎Broadcast┓ + | | | + 2 3 4 +``` + +*Four branches, so four copies of each watermark must arrive before its offsets are +committed.* + +Because a watermark is committed only once every copy has arrived, its offsets correspond +to work finished on every path. Because watermarks are periodic rather than per-message, +the bookkeeping is bounded: the commit step tracks a handful of in-flight watermarks, not a +set of in-flight offsets. Trackers that never receive all their copies — a branch that +broke and stopped forwarding — are dropped after a timeout so the buffer cannot grow +without bound. + +Most steps forward watermarks. The ones that do not are the steps that hold messages back: + +| Step | Watermark handling | +| --- | --- | +| **Reduce** | Accumulates the committables of watermarks received while a window is open; releases the combined result when the window closes | +| **Python delegate** | Holds a watermark until the output produced covers the offsets it carries | +| **Broadcast / Router** | Duplicate them, one per branch | + +Watermarks also carry the timestamp of the newest data message seen since the previous +watermark, which is how end-to-end consumer latency is measured at commit time. + +### Crossing the language boundary + +A watermark is a `Watermark` struct in Rust. When it must be handed to Python — a delegate +needs to see it to keep the ordering guarantee above — it is converted into `PyWatermark`, a +`pyclass` whose committable is a Python dict, and converted back on the way out. That is +why `WatermarkMessage` has two variants, and it comes with directional invariants recorded +in [Contracts](../contracts.md#the-message-model). diff --git a/sentry_streams/docs/architecture/rust-arroyo/metrics.md b/sentry_streams/docs/architecture/rust-arroyo/metrics.md new file mode 100644 index 00000000..0ff2859b --- /dev/null +++ b/sentry_streams/docs/architecture/rust-arroyo/metrics.md @@ -0,0 +1,61 @@ +# Metrics + +**Audience:** everyone · **Durability:** mechanism · **Last reviewed:** 2026-09-15 + +A running pipeline emits metrics from both sides of the language boundary and from several +independent producers within each. They are configured once and all end up in the same +place under the same namespace. + +How to configure them is in +[Deployment configuration](../../reference/deployment-config.md#metrics). This page covers +only what is architecturally load-bearing. + +## What produces metrics + +| Producer | Side | What it reports | +| --- | --- | --- | +| Pipeline stats | Python | Per-step message counts, errors, durations, for steps executed in Python | +| Pipeline stats | Rust | The same, for steps executed natively | +| Rust strategies | Rust | Operator internals: Python adapter submit/poll durations, batch sizes, commit latency | +| Arroyo (Python) | Python | The Python Arroyo library's own consumer metrics | +| Arroyo (Rust) | Rust | The Rust Arroyo library's own consumer metrics | + +All are namespaced under `streams.pipeline`. Tags carry at least the pipeline name, and +per-step metrics carry the step name, so a dashboard can break a pipeline down by step +regardless of which side of the boundary the step ran on. + +## Deliberate symmetry across the boundary + +Step-level metrics are handled by an equivalent buffering component on each side, with the +same semantics and the same metric names: + +- In Python, `PipelineStats` accumulates counters and timings per step and flushes every ten + seconds. The adapter wraps every application function it registers, so counts and timings + are recorded around the call, including when it raises. +- In Rust, the same buffering exists with thread-local state (`src/pipeline_stats.rs`), + flushed on the same interval, emitting the same names with the same `step` tag. + +The point of the symmetry is that a step reports identically whether it was fused into a +native operator or executed as a Python function, so one set of dashboards reads a pipeline +end to end. A new Rust strategy is expected to keep that property — see +[Adding a step](../guides/adding-a-step.md#6-instrument-it). + +Buffering rather than sampling is deliberate: some metrics are produced in tight loops where +emitting costs a significant fraction of the work being measured, and aggregation keeps rare +events instead of discarding them. + +Beyond step stats, the Rust strategies emit their own internals, which is where to look when +a pipeline is slow and step timings do not explain it: how long a `submit` into a Python +delegate took, how long the delegate's `poll` took, how long the next strategy took to +accept the result, and end-to-end consumer latency measured at commit time from the +timestamp a watermark carries. + +## Multiprocessing + +Worker processes are separate interpreters and inherit nothing. The pool is created with an +initializer that calls `configure_metrics` again, with the same config, in each worker. + +This is why the metrics configuration is a plain dictionary rather than a live backend +object: it has to survive being pickled and sent to a worker. The broader rule — config is +data, re-applied per process, never passed by reference — is an +[invariant](../contracts.md#configuration). diff --git a/sentry_streams/docs/architecture/rust-arroyo/python-operator.md b/sentry_streams/docs/architecture/rust-arroyo/python-operator.md new file mode 100644 index 00000000..31034d83 --- /dev/null +++ b/sentry_streams/docs/architecture/rust-arroyo/python-operator.md @@ -0,0 +1,139 @@ +# The Python operator + +Most Python in this runtime is a function called from a Rust strategy: the map, the filter, +the routing function. The Python operator is for the cases where that is not enough — where +the thing that has to run in Python is not a function but a *step*, with its own buffering, +its own timing, and its own output cadence. + +Two such steps exist today, both pre-existing Python Arroyo strategies that were not +reimplemented in Rust: the **multiprocessing** step and the **windowed reduce**. + +## The delegate interface + +`RuntimeOperator::PythonAdapter` builds `PythonAdapter` (`src/python_operator.rs`), a Rust +Arroyo strategy holding a reference to a Python object and delegating message processing to +it. To the rest of the chain it is an ordinary `ProcessingStrategy`: it can be +wired between any two strategies, it propagates backpressure, it participates in commits. + +What it delegates to is not an Arroyo strategy, and that is the crux of the design. An +Arroyo strategy hands its results to the next strategy itself, which cannot work across the +boundary: the next strategy is a Rust object and a Python object cannot hold or call it. So +the Python side implements a deliberately simpler interface, `RustOperatorDelegate`: + +| Method | Contract | +| --- | --- | +| `submit(message, committable)` | Accept work. It does not process, it stores | +| `poll()` | Do the processing and **return** the results, as `(message, committable)` pairs | +| `flush(timeout)` | Finish everything in flight, return the results, release resources | + +Rust forwards whatever `poll` returned to the next strategy, so the delegate never needs a +reference to anything downstream. + +Splitting "accept" from "produce" also makes the interface inherently asynchronous: a +delegate can be 1:1, 1:N, N:1 or N:0 — a reduce accepts a hundred messages and produces +one, a multiprocess pool accepts a batch and produces results several polls later. None of +that is expressible in a signature that must return one result per input. + +A delegate is created by a `RustOperatorFactory`, not passed in directly, because Arroyo +tears down and rebuilds its strategy chain on every rebalance. The factory is also where +state that must survive a rebuild lives — a pre-initialised process pool, for instance, +far too expensive to recreate on every partition assignment. + +## Data flow + +```mermaid +sequenceDiagram + participant Prev as previous Rust strategy + participant PA as PythonAdapter (Rust) + participant Del as delegate (Python) + participant Next as next Rust strategy + + Prev->>PA: submit(Message) + Note over PA: route check — mismatch forwards straight to Next + Note over PA: extract committable + PA->>PA: acquire GIL + PA->>Del: submit(py_payload, py_committable) + Note over Del: stores the work + Del-->>PA: None / MessageRejected / InvalidMessage + + loop processor main loop + PA->>Del: poll() + Del-->>PA: [(payload, committable), ...] + PA->>PA: wrap each into Message + PA->>Next: submit(...) for each, in order + end +``` + +Properties of that flow which are not obvious from the diagram: + +- The GIL is acquired **once** per `submit`, and both conversions happen inside it: payload + to a Python object, committable to a dict keyed by `(topic, partition)`. A `PyAnyMessage` + or `RawMessage` is already a `pyclass`, so this is a reference clone, not a copy. +- Exceptions out of `submit` are the delegate's control channel: + + | Exception | Meaning | Effect | + | --- | --- | --- | + | `MessageRejected` | "I am full" | Arroyo backpressure; the message is handed back to be retried | + | `InvalidMessage` | "this cannot be processed" | Offset and partition are read off the exception, message goes to the DLQ | + | anything else | a bug in the delegate | panic — there is no sensible interpretation, and continuing would silently drop data | + +- On the way out, the payload decides the variant: a `PyWatermark` becomes a watermark + message, anything else a `PyStreamingMessage`. **The route attached is the operator's + own** — the delegate has no notion of routes and does not need one. +- Results are drained into the next strategy one at a time, polling it between + submissions. If it rejects one, that message goes back to the front of the queue and the + drain stops, preserving order and propagating backpressure upstream. `join` follows the + same path through `flush`, with a deadline. +- The committable a delegate returns is *its* choice, not an echo of the input: a reduce + that merged a hundred messages returns their combined committable. The messages produced + are Arroyo "any messages" rather than broker messages, which is why a failure downstream + of this operator cannot be attributed to a specific offset. + +## Wrapping a whole Arroyo strategy + +The delegate interface is small enough to implement directly — `SingleMessageOperatorDelegate` +is a helper for the trivial 1:1 synchronous case. But both real users wrap an existing +Python Arroyo strategy, unmodified, via `ArroyoStrategyDelegate`, which bridges three +mismatches: + +| Mismatch | Bridge | +| --- | --- | +| The strategy pushes; the delegate returns | `OutputRetriever`, a minimal Arroyo strategy given to the wrapped strategy as its next step. It forwards nowhere and collects what it receives into a list, which `poll` hands back | +| The strategy speaks Arroyo messages; the runtime speaks pipeline messages | Two transformer functions, in and out. The output transformer is also where Arroyo-specific payloads like `FilteredPayload` are dropped | +| Watermarks must not overtake the data | The delegate intercepts watermarks, holds them, and releases one only once the output produced covers the offsets it carries | + +```mermaid +flowchart LR + rust["PythonAdapter (Rust)"] + del["ArroyoStrategyDelegate"] + inner["Wrapped Arroyo strategy
(multiprocess / reduce)"] + ret["OutputRetriever"] + + rust -- "submit(payload, committable)" --> del + del -- "in_transformer → ArroyoMessage" --> inner + inner -- "submit (its next_step)" --> ret + ret -- "out_transformer" --> del + del -- "poll() → results" --> rust +``` + +## The multiprocess step + +The main user of this machinery is `RunTaskWithMultiprocessing`, Arroyo's pool-based +transformation strategy. The adapter reaches for it when a segment is configured with +`multi_process` parallelism: the fused chain of maps becomes the function executed in the +pool, and the whole strategy is wrapped in a delegate. + +One constraint shapes the conversion. `RunTaskWithMultiprocessing` moves payloads to worker +processes through shared memory, which means pickling them — and the Rust message types are +`pyclass`es, which are not picklable. + +That is why the pipeline messages Python code sees are the `Message` wrappers (`PyMessage`, +`PyRawMessage`, in `sentry_streams/pipeline/message.py`) rather than the Rust types +directly. The wrappers hold plain Python fields, pickle cleanly, and build the underlying +Rust object lazily on first use, caching it. The input transformer unwraps a Rust message +into a wrapper; the output transformer calls `to_inner()` to get the Rust object back. The +same wrappers are what make the payload generic to the type checker, which `pyclass` types +cannot be. + +Worker processes are a fresh interpreter each, so the pool is created with an initializer +that reconfigures metrics in every worker — see [Metrics](./metrics.md#multiprocessing). diff --git a/sentry_streams/docs/reference/deployment-config.md b/sentry_streams/docs/reference/deployment-config.md new file mode 100644 index 00000000..160d38ac --- /dev/null +++ b/sentry_streams/docs/reference/deployment-config.md @@ -0,0 +1,184 @@ +# Deployment configuration + +The deployment config is a YAML file, validated against a JSON schema, holding everything +that varies between deployments of the same application. Its purpose is to keep +infrastructure concerns out of product code: the same application file, deployed with a +different config, reads from a different topic, runs with different parallelism, and reports +to a different metrics backend. + +This is the single description of the file. Other documents link here rather than repeating +it. + +- Schema: `sentry_streams/config.json` +- Types: `sentry_streams/deployment_config/config_types.py` +- Examples: `sentry_streams/deployment_config/*.yaml` + +## Structure + +```yaml +env: {} # general environment settings +metrics: # metrics backend, both languages + type: log + period_sec: 5 + tags: + pipeline: errors +sentry_sdk_config: # error reporting + dsn: "${envvar:SENTRY_DSN}" +pipeline: + adapter_config: # settings meaningful to one runtime only + arroyo: + write_healthcheck: true + segments: # deployment units; --segment-id selects one + - steps_config: # keyed by step name + myinput: + starts_segment: True + bootstrap_servers: ["127.0.0.1:9092"] + parser: + starts_segment: True + parallelism: + multi_process: + processes: 4 + batch_size: 1000 + batch_time: 0.2 + mysink: + starts_segment: True + bootstrap_servers: ["127.0.0.1:9092"] +``` + +## Environment variables + +Any value may reference an environment variable with a `${envvar:NAME}` placeholder, resolved +when the config is loaded, so secrets and per-environment endpoints need not live in the file. + +```yaml +override_params: + max.poll.interval.ms: "${envvar:MAX_POLL_INTERVAL_MS}" +``` + +## Segments + +The word has two meanings. Both are in the [glossary](../architecture/README.md#glossary); +this is how each appears in the file. + +**`pipeline.segments` is a list.** Each entry is a deployment unit with its own `steps_config`, +selected at run time with `--segment-id`. The adapter only ever sees the configuration of the +segment it is running. The Kubernetes integration deploys one workload per entry. + +**`starts_segment: True` inside a `steps_config`** marks a boundary *within* what a single +process runs: it is where the adapter stops fusing consecutive maps and where a new +parallelism setting takes effect. + +Because the Rust adapter builds a single consumer per process, the list usually has one entry +and the interesting segmentation is the one done by `starts_segment`. + +### Parallelism + +Declared on the step that starts a segment; configuring it mid-chain is rejected. + +```yaml +parser: + starts_segment: True + parallelism: + multi_process: + processes: 4 + batch_size: 1000 + batch_time: 0.2 +``` + +The fused function of the segment becomes the function executed in the pool, which is why it +must be picklable — a module-level function pickles, a closure or local function does not, and +the check runs at build time even when multiprocessing is off. + +## Per-step configuration + +`steps_config` is keyed by step name. Keys fall into two groups, which is worth knowing when a +value appears to be ignored: + +**Read by the step**, in its `override_config()`: + +| Step | Keys | +| --- | --- | +| `StreamSource` | `topic`, `consumer_group` | +| `StreamSink` | `topic` | +| `GCSSink` | `bucket`, `parallelism.threads` | +| `Batch` | `batch_size`, `batch_timedelta` (a mapping of `timedelta` kwargs) | +| `DevNullSink` | `batch_size`, `batch_time_ms`, `average_sleep_time_ms`, `max_sleep_time_ms` | + +**Read by the adapter**, not by the step: `starts_segment`, `parallelism`, +`bootstrap_servers`, `override_params` (passed through to the Kafka client), and `dlq`. + +A source keeps the **logical** stream name written in the application for schema and codec +lookup even when `topic` is overridden: the logical name identifies the data, the topic +identifies where it lives in this environment. + +Adding a new key for a new step needs **no schema change** — `steps_config` entries allow +additional properties. Document the keys in the step's docstring and add a row above. + +## Metrics + +```yaml +metrics: + type: datadog + host: 127.0.0.1 + port: 8125 + tags: + environment: production + flush_interval_ms: 1000 +``` + +| Type | Behaviour | +| --- | --- | +| `datadog` | Sends to a DogStatsD agent over UDP | +| `log` | Writes each metric to the log, with a configurable `period_sec`. Useful locally | +| `dummy` | Discards. The default when no `metrics` block is present, and what tests use | + +The runner adds the pipeline name as a default tag and configures both languages from this one +block: the Python backend plus an adapter installed into the Python Arroyo library, and — via +`PyMetricConfig` — a DogStatsD exporter and a recorder for Rust Arroyo. Note the asymmetry: +the Rust side only has a real backend for `datadog`, so with `log` or `dummy` the Rust metrics +are not produced. + +What the metrics mean is in [Metrics](../architecture/rust-arroyo/metrics.md). + +## Adapter configuration + +`pipeline.adapter_config` holds settings meaningful to one runtime only. For the Arroyo +adapters: + +| Key | Effect | +| --- | --- | +| `arroyo.write_healthcheck` | Adds the `HealthCheck` strategy, which touches a file on every poll for the Kubernetes liveness probe | +| `arroyo.sentry_sdk_config` | Sentry SDK settings for the adapter | + +## Dead letter queue + +A `dlq` block on a source configures the DLQ for that consumer: `topic`, `bootstrap_servers`, +`override_params`. **With no DLQ configured, an invalid message stops processing** — the +consumer logs this at startup. Current DLQ limits are Arroyo's defaults (no rate limit, no cap +on buffered messages); + +## The seam with `sentry_streams_k8s` + +The `sentry_streams_k8s` package renders the Kubernetes objects that run a pipeline. The +contract between the two packages is this file: + +- One **workload per `pipeline.segments` entry**, each started with the matching + `--segment-id`. +- The config file itself is delivered as a **ConfigMap**, rendered by the sentry-kube consumer + macro (or by the experimental operator) alongside the Deployment. +- `arroyo.write_healthcheck` is what makes the liveness probe meaningful, so it and the probe + have to be enabled together. + +See `sentry_streams_k8s/README.md` for the rendering side. + +## Example files + +| File | Shows | +| --- | --- | +| `simple_map_filter.yaml` | The minimum: source and sink bootstrap servers | +| `parallel_processing.yaml` | Segments and `multi_process` parallelism | +| `envvars.yaml` | `${envvar:...}` placeholders and `override_params` | +| `simple_batching.yaml` | `Batch` configuration | +| `gcs_sink.yaml` | GCS sink configuration | +| `devnull_benchmark.yaml` | Benchmarking with `DevNullSink` | +| `blq.yaml` | The branching `blq.py` example: broadcast, router, several Kafka sinks | diff --git a/sentry_streams/sentry_streams/deployment_config/README.md b/sentry_streams/sentry_streams/deployment_config/README.md index b389c022..e66993e9 100644 --- a/sentry_streams/sentry_streams/deployment_config/README.md +++ b/sentry_streams/sentry_streams/deployment_config/README.md @@ -1,31 +1,10 @@ -# Configuration files +# Example configuration files -Configuration files have two main sections right now; `env` and `pipeline`. These files are also -adapter/runtime-specific right now. +This directory holds example deployment configuration files, the JSON schema +(`../config.json`) and the concrete types (`config_types.py`). -`env` is supposed to be where all general environment config can go. For example, in the -case of a Flink configuration file, that could mean setting certain properties that hold true -for an entire streaming pipeline. - -`pipeline` holds optional runtime-level options and configuration for each segment. -You can set `pipeline.adapter_config.arroyo.write_healthcheck: true` to enable -the Arroyo healthcheck strategy (touches a file for Kubernetes liveness probes). A segment is defined as a set of steps or operators -which will be executed on one unit (this could be a worker, a single consumer, or a Flink slot). Each -segment can have a parallelism override if parallelism is a value other than 1. The segment also holds -step-specific configuration as a mapping. This step config should hold values that are overrides of defaults. -As of now, the only steps that need specific configuration are sources and sinks. Thus, segments which contain -either a source or sink must have the config and any overrides specified. - -The idea is that, using segments and parallelism, we can achieve distribution of steps across different -physical workers. With a runtime like Flink, we can create segments (also known as chains in Flink terms) -for different parts of the pipeline, and give each segment different parallelism values. With this, we -could have, for example, a segment which lives in one worker, reshuffling data to another segment which -is distributed across 10 workers. - -For now, step-specific defaults (for settings like `auto-offset-reset`) are embedded into each adapter, -but ultimately these should live in configuration. - -See `config.json` for the current schema of a configuration file and `config_types.py` for concrete types -that the configuration gets converted into. - -See example configuration files in this directory. +**The configuration file format is documented once, in +[docs/reference/deployment-config.md](../../docs/reference/deployment-config.md).** That page +describes the structure, segments and parallelism, per-step keys, environment variable +placeholders, metrics, the DLQ, and the seam with the Kubernetes integration — and lists what +each example in this directory demonstrates. From 7083a27e394c288780b741dad9a51c05ea9cbfc6 Mon Sep 17 00:00:00 2001 From: Filippo Pacifici Date: Wed, 16 Sep 2026 09:44:14 -0700 Subject: [PATCH 2/3] Limit arroyo --- sentry_streams/Cargo.toml | 2 +- sentry_streams/pyproject.toml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/sentry_streams/Cargo.toml b/sentry_streams/Cargo.toml index 07847758..e7780d55 100644 --- a/sentry_streams/Cargo.toml +++ b/sentry_streams/Cargo.toml @@ -6,7 +6,7 @@ edition = "2021" [dependencies] pyo3 = { version = "0.29.0" } serde = { version = "1.0", features = ["derive"] } -sentry_arroyo = { version = "2.40.0", features = ["ssl"] } +sentry_arroyo = { version = ">=2.40.0, <2.44.0", features = ["ssl"] } chrono = "0.4.40" tracing = "0.1.40" tracing-subscriber = "0.3.20" diff --git a/sentry_streams/pyproject.toml b/sentry_streams/pyproject.toml index d5c3e06d..cae60af6 100644 --- a/sentry_streams/pyproject.toml +++ b/sentry_streams/pyproject.toml @@ -18,7 +18,7 @@ classifiers = [ dependencies = [ "requests>=2.32.3", - "sentry-arroyo>=2.39.2", + "sentry-arroyo>=2.39.2,<2.44.0", "pyyaml>=6.0.2", "jsonschema>=4.23.0", "sentry-kafka-schemas>=2.2.0", From 87c12eb370459b8cdf308f52f5725fbc5033768b Mon Sep 17 00:00:00 2001 From: Filippo Pacifici Date: Wed, 16 Sep 2026 11:52:20 -0700 Subject: [PATCH 3/3] Add decisions --- sentry_streams/docs/architecture/README.md | 6 +- ...01-separate-application-from-deployment.md | 49 +++++++++++++ .../0002-adapters-as-the-runtime-boundary.md | 37 ++++++++++ .../decisions/0004-complex-steps-as-sugar.md | 46 ++++++++++++ .../decisions/0005-statically-typed-dsl.md | 35 ++++++++++ .../decisions/0009-rust-runs-python-builds.md | 70 +++++++++++++++++++ .../decisions/0010-routes-for-branching.md | 46 ++++++++++++ .../0011-watermarks-drive-commits.md | 50 +++++++++++++ ...d-python-strategies-via-a-pull-delegate.md | 44 ++++++++++++ .../decisions/0014-fuse-maps-into-segments.md | 49 +++++++++++++ .../0018-symmetric-buffered-metrics.md | 45 ++++++++++++ .../docs/architecture/decisions/README.md | 43 ++++++++++++ 12 files changed, 518 insertions(+), 2 deletions(-) create mode 100644 sentry_streams/docs/architecture/decisions/0001-separate-application-from-deployment.md create mode 100644 sentry_streams/docs/architecture/decisions/0002-adapters-as-the-runtime-boundary.md create mode 100644 sentry_streams/docs/architecture/decisions/0004-complex-steps-as-sugar.md create mode 100644 sentry_streams/docs/architecture/decisions/0005-statically-typed-dsl.md create mode 100644 sentry_streams/docs/architecture/decisions/0009-rust-runs-python-builds.md create mode 100644 sentry_streams/docs/architecture/decisions/0010-routes-for-branching.md create mode 100644 sentry_streams/docs/architecture/decisions/0011-watermarks-drive-commits.md create mode 100644 sentry_streams/docs/architecture/decisions/0013-embed-python-strategies-via-a-pull-delegate.md create mode 100644 sentry_streams/docs/architecture/decisions/0014-fuse-maps-into-segments.md create mode 100644 sentry_streams/docs/architecture/decisions/0018-symmetric-buffered-metrics.md create mode 100644 sentry_streams/docs/architecture/decisions/README.md diff --git a/sentry_streams/docs/architecture/README.md b/sentry_streams/docs/architecture/README.md index d73cc4f3..49b17ea3 100644 --- a/sentry_streams/docs/architecture/README.md +++ b/sentry_streams/docs/architecture/README.md @@ -2,7 +2,8 @@ This directory describes the architecture of the `sentry_streams` package: the roles of the main subsystems, the decisions behind them, the rules that must hold, and what we -know is unfinished. +know is unfinished. The decisions themselves, one per page, are in the +[decision log](./decisions/README.md). # The model @@ -84,7 +85,7 @@ The DSL is built from a small set of primitives, in a few families: ones: `Parser`, `Serializer`, `BatchParser`, `ParquetSerializer`, `Reducer`. See [Pipeline DSL and runner](./pipeline-dsl-and-runner.md#complex-steps). -Refer to [pipeline.py](../../sentry_streams/pipeline/pipeline.py) to find documentaiton +Refer to [pipeline.py](../../sentry_streams/pipeline/pipeline.py) to find documentation on each of the implemented steps. A pipeline is built by chaining these together. Each step has a name, which is also the @@ -158,6 +159,7 @@ in [The Rust Arroyo adapter](./rust-arroyo/README.md). **Architecture** +1. [Decision log](./decisions/README.md) — the choices this design rests on, and why. 2. [Guarantees and failure modes](./guarantees.md) — what the system promises. diff --git a/sentry_streams/docs/architecture/decisions/0001-separate-application-from-deployment.md b/sentry_streams/docs/architecture/decisions/0001-separate-application-from-deployment.md new file mode 100644 index 00000000..7c4d9ede --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0001-separate-application-from-deployment.md @@ -0,0 +1,49 @@ +# 0001 — An application describes a dataflow; a separate config describes its deployment + + +## Decision + +A streaming application is a pure, declarative description of a dataflow graph: a source, a +series of transformations, one or more sinks. It contains no broker address, no consumer +group, no topic, no parallelism setting and no reference to the runtime that will execute +it. Everything the description leaves out lives in a separate deployment configuration +file, which addresses each step by its name. + +The split is a hard line, not a convention: only non-semantic values are overridable +(see [Configuration](../contracts.md#configuration)), and step names are the only coupling +between the two artefacts. + + +## Rationale + +The decision is meant to simplify the separation between infrastructure and application +which is needed to build a streaming platform. + +Everything in the config file is specific to how the pipeline is operated: + +* Different by environment (Kafka hosts) + +* Scale configuration + +* Optimizations + +Separating these concepts allows the product engineers to focus on how the +application works while platform engineers can ensure the application runs in +a scalable and performant way. + +Enforcing a separation ensures also that platform engineers would have the same +parameters to tune for all pipelines independently on what the application +does. + +## Consequences + +- The same application file can be deployed against different brokers, topics, parallelism + settings and even different runtimes without being modified. +- Renaming a step is a breaking change to its deployment configuration. +- Anything a deployment needs to vary must be modelled as an overridable key on a step; + there is no escape hatch. + +## See also + +[The model](../README.md#the-model) · [Pipeline DSL and runner](../pipeline-dsl-and-runner.md) · +[Deployment configuration](../../reference/deployment-config.md) diff --git a/sentry_streams/docs/architecture/decisions/0002-adapters-as-the-runtime-boundary.md b/sentry_streams/docs/architecture/decisions/0002-adapters-as-the-runtime-boundary.md new file mode 100644 index 00000000..e1835ab6 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0002-adapters-as-the-runtime-boundary.md @@ -0,0 +1,37 @@ +# 0002 — Runtimes are pluggable adapters with one method per primitive + +## Decision + +The execution runtime sits behind an adapter interface with one abstract method per +primitive (`source`, `map`, `filter`, `reduce`, `sink`, `router`, `broadcast`, `flat_map`). +The interface is generic in the type of the stream handle, so an adapter represents "a +stream" however its runtime requires, and the runner never inspects it. Dispatch goes +through `RuntimeTranslator` over a closed `StepType` enum ending in `assert_never`. + +Adding a primitive is therefore a **breaking change to every adapter**, by construction. An +adapter that will not support a primitive declines explicitly with `NotImplementedError`, +as the Rust adapter does for `flat_map`. + + +## Rationale + +At the time of writing the team is not settled on which runtime should be used +at scale. This abstraction layer allows the streaming team to experiment with +different runtimes while not breaking existing applications. + +Having one method per primitive is meant to reduce the complexity of the +adapter. + +## Consequences + +- A new primitive cannot be added silently: the type checker fails until every adapter and + the translator handle it. +- Pressure toward [complex steps](./0004-complex-steps-as-sugar.md), which cost no adapter + changes. +- The runner stays thin; runtime knowledge lives entirely in the adapter. + +## See also + +[Extending the DSL](../contracts.md#extending-the-dsl) · +[Adding a step](../guides/adding-a-step.md) · +[From pipeline to consumer](../pipeline-dsl-and-runner.md#from-pipeline-to-consumer) diff --git a/sentry_streams/docs/architecture/decisions/0004-complex-steps-as-sugar.md b/sentry_streams/docs/architecture/decisions/0004-complex-steps-as-sugar.md new file mode 100644 index 00000000..29532462 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0004-complex-steps-as-sugar.md @@ -0,0 +1,46 @@ +# 0004 — Higher level steps are sugar that converts to primitives, with a native override hook + +## Decision + +Recurring combinations of primitives — `Parser`, `Serializer`, `BatchParser`, +`ParquetSerializer`, `Reducer` — are `ComplexStep`s. Each implements `convert()`, returning +plain simple steps, usually a `Map` with a function partial-applied. The translator calls +`convert()` transparently, so by the time an adapter sees the pipeline a `Parser` is just a +`Map`. + +An adapter that has a faster native implementation of a complex step declares it through +`complex_step_override()` and receives the original step instead of its conversion. +Adapters that do not care return an empty mapping. + +## Alternatives + +- Make every useful step a primitive, and require every adapter to implement it. +- Keep them as plain helper functions in application code. + +## Rationale + +There are many common variants of basic steps that customers would have to implement +on their own starting from basic building block. + +If we made a primitive for each of them we would have to implement them on each +runtime in depth. + +If we asked product teams to implement each we would have a lot of cases where +we reinvent the wheel. + +This allows us to provide composite steps quickly without the overhead of +implementing them in each adapter. + + +## Consequences + +- The set of operations an adapter must implement stays small. +- Complex steps are the first thing to reach for when extending the DSL; only what cannot be + composed becomes a primitive. +- A complex step's per-step configuration and metrics are attributed to the step it converts + into, unless an adapter overrides it natively. + +## See also + +[Complex steps](../pipeline-dsl-and-runner.md#complex-steps) · +[Decide what you are actually adding](../guides/adding-a-step.md#0-decide-what-you-are-actually-adding) diff --git a/sentry_streams/docs/architecture/decisions/0005-statically-typed-dsl.md b/sentry_streams/docs/architecture/decisions/0005-statically-typed-dsl.md new file mode 100644 index 00000000..f0ec6dc8 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0005-statically-typed-dsl.md @@ -0,0 +1,35 @@ +# 0005 — The DSL is statically typed end to end + +## Decision + +`Pipeline[TOut]` carries the type of the messages flowing out of the last step added, and +every primitive is generic in its input and output types. Chaining a step whose input type +does not match the pipeline's current output type is a type error, caught by mypy in strict +mode rather than at runtime. + +## Alternatives + +- Untyped payloads, with mismatches surfacing as runtime failures in the consumer. +- Runtime schema checks at step boundaries instead of static types. + +## Rationale + +This is meant to reduce the chances of runtime errors though it has a cost: +errors are not very easy to spot. + +We should still consider validating the types at startup or during tests. + +## Consequences + +- Structural mistakes are caught before deployment — for instance sinking parsed objects + into a Kafka sink with no serializer in between, which would otherwise panic at runtime + (see [Messages](../rust-arroyo/messages.md#conversion-is-limited-on-purpose)). +- Every new step class must declare its generic parameters correctly for the chain to keep + type checking. +- Rust `pyclass` types cannot be made generic to the type checker, which is one reason + Python code sees the `Message` wrappers rather than the Rust types. + +## See also + +[Principles](../pipeline-dsl-and-runner.md#principles) · +[Declare the step in the DSL](../guides/adding-a-step.md#2-declare-the-step-in-the-dsl) diff --git a/sentry_streams/docs/architecture/decisions/0009-rust-runs-python-builds.md b/sentry_streams/docs/architecture/decisions/0009-rust-runs-python-builds.md new file mode 100644 index 00000000..9962f4c4 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0009-rust-runs-python-builds.md @@ -0,0 +1,70 @@ +# 0009 — Rust owns the runtime loop; Python builds the pipeline + +## Decision + +The production runtime is an Arroyo consumer written in Rust, exposed to Python as a PyO3 +extension module. The process starts as Python, imports the extension, and hands control +over: **Python builds** a description of the pipeline and passes it to Rust operator by +operator; **Rust runs**, owning the main loop, polling Kafka, submitting through the +strategy chain and committing; **Rust calls back into Python** for application logic. + +The governing rule for what goes where is that **any primitive that can be executed +entirely in Rust is implemented in Rust**, because crossing the boundary costs a GIL +acquisition, serialises the step against every other Python interaction in the process, and +may cost a copy. + +## Alternatives + +- A pure Python runtime on Python Arroyo (the pre-existing adapter, still in the tree). +- A pure Rust process with no embedded interpreter, requiring application logic in Rust. +- A Rust runtime that starts the Python interpreter and calls into Python + for Python logic. + +## Rationale + +This platform is architected to support Python applications, Rust applications +and hybrid ones. Managing the main loop in Rust allows us to get the best +performance specifically from Rust applications. + +This allows us to support these scenarios: + +* For pure Rust applications data would never move into Python memory and the + GIL is never used. This is not fully implemented though. + +* There can be two types of Python hybrid applications, one where each message + moves to Python memory and the other where the logic is written in Python + but the processing happens in Rust. + +The last scenario is particularly important, it can be implemented this way: + +- Use a declarative query language to define the application logic in Python + with DataFusion + +- The processing, though, happens fully in Rust and is done by DataFusion. + +Running the loop in Rust allows us not to move data to Python for DataFusion +processing. + +Making the runner a Python process is done for expedience though. We built the +DSL and the first runtime in Python, most applications are written in Python +or are in Python code bases. + +Having a Python DSL is desirable to support those. We could still have a Rust +runner to replace the Python one (which will be desirable), it is just not implemented +yet. + +## Consequences + +- Kafka consumption, offset management, the strategy chain, backpressure, DLQ routing, + watermarks, route dispatch, batching, header filtering, broadcast and all sinks are Rust. +- Application maps, filters and routing functions, parsing and serialisation, windowed + reduce and multiprocessing remain Python. +- Every design question in this directory — the message model, routes, the delegate + interface — exists because of this boundary. +- A middle ground is available: a function written in Rust, exposed as a Python callable, + releasing the GIL internally. + +## See also + +[Execution model](../rust-arroyo/README.md#execution-model) · +[What runs where](../rust-arroyo/README.md#what-runs-where) diff --git a/sentry_streams/docs/architecture/decisions/0010-routes-for-branching.md b/sentry_streams/docs/architecture/decisions/0010-routes-for-branching.md new file mode 100644 index 00000000..86e94747 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0010-routes-for-branching.md @@ -0,0 +1,46 @@ +# 0010 — Branching is expressed as routes on a single linear strategy chain + +**Status:** accepted · **Area:** Rust runtime + +## Decision + +Arroyo pipelines are linear: a strategy has exactly one next strategy. The DSL has `Router` +and `Broadcast`. Rather than building a tree of strategies, every message carries a +`Route` — its source name plus the ordered waypoints of the branches it has been sent +down — and every strategy is built with the route it belongs to and forwards untouched +anything whose route differs. The chain stays physically linear and holds the steps of +every branch one after another. + +Consequently the Rust adapter's stream handle **is** a `Route`: there is no stream object to +hand around, so what the runner passes to an adapter method is "which branch", and what it +gets back is the branch the output belongs to — one route per branch for router and +broadcast. + +## Alternatives + +- Extend Arroyo with real branching strategies that own several next steps. +- One consumer process per branch. + +## Rationale + +Keeping a single arroyo chain of steps simplified considerably commit management +to guarantee at least once delivery. + +Having one single chain ensures that there is a single last step that receives +messages from all branches and can evaluate which offsets have to be committed. + +## Consequences + +- A message physically traverses every strategy in the consumer, including branches it will + never enter. This is affordable only because the route check is cheap and GIL-free, which + makes "check the route first, forward on mismatch, take no GIL" the first rule for every + step author. +- Commit correctness no longer follows from reaching the end of the chain, which is what + [0011](./0011-watermarks-drive-commits.md) exists to fix. +- More than one consumer per process is not supported today, though the adapter's structure + anticipates it. + +## See also + +[Routes](../rust-arroyo/messages.md#routes) · +[The stream handle is a route](../rust-arroyo/adapter.md#the-stream-handle-is-a-route) diff --git a/sentry_streams/docs/architecture/decisions/0011-watermarks-drive-commits.md b/sentry_streams/docs/architecture/decisions/0011-watermarks-drive-commits.md new file mode 100644 index 00000000..badaf74f --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0011-watermarks-drive-commits.md @@ -0,0 +1,50 @@ +# 0011 — Commits are driven by periodic watermarks, not by data messages + +## Decision + +A `WatermarkEmitter` at the head of the chain records the offsets of every message passing +through it and, every few seconds (10 by default), injects a watermark control message +carrying the accumulated committable. Watermarks travel down the pipeline like data +messages and go through the same branching logic. At the end of the chain the commit policy +ignores data messages entirely: it counts watermark copies, and commits the carried offsets +only once it has seen that watermark arrive from every branch. + +Steps that hold messages back hold watermarks with them; every other step forwards them. + +## Alternatives + +- Arroyo's default policy: commit a message's offsets when it reaches the end of the chain. +- Track every in-flight offset and its per-branch completion. + +## Rationale + +Watermarks make it much easier to calculate the offset to commit and when to commit +in scenarios where the pipeline has branches. + +The last step of the pipeline receives all the watermarks propagated through all +branches and is able to ensure that all branches reached a specific offset +before committing. + +This also allows us to drop messages (via filters for example) and it allows each +step to decide when it is done processing up a certain offsets just by holding +the offsets instead of propagating them. + +## Consequences + +- With a broadcast in the pipeline, committing on the first copy to arrive would commit work + that has not happened; counting copies is what makes at-least-once hold under branching. +- The bookkeeping is bounded: a handful of in-flight watermarks rather than a set of + in-flight offsets. Trackers that never receive all their copies are dropped after a + timeout. +- Commit latency is bounded below by the watermark interval, not by message throughput. +- Every step author must handle watermarks explicitly, and buffering steps must not let them + overtake the data. +- A watermark crossing into Python becomes a `PyWatermark`, with directional invariants: + never submitted back into a Python operator, never handed to the Kafka sink. +- Watermarks carry the newest data message timestamp, which is how end-to-end consumer + latency is measured at commit time. + +## See also + +[Watermarks](../rust-arroyo/messages.md#watermarks) · +[The message model](../contracts.md#the-message-model) diff --git a/sentry_streams/docs/architecture/decisions/0013-embed-python-strategies-via-a-pull-delegate.md b/sentry_streams/docs/architecture/decisions/0013-embed-python-strategies-via-a-pull-delegate.md new file mode 100644 index 00000000..fa7a6649 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0013-embed-python-strategies-via-a-pull-delegate.md @@ -0,0 +1,44 @@ +# 0013 — Existing Python Arroyo strategies are embedded through a pull-based delegate, not reimplemented + +## Decision + +Where the thing that must run in Python is not a function but a *step* — with its own +buffering, timing and output cadence — it runs behind `RuntimeOperator::PythonAdapter`, a +Rust strategy that delegates to a Python object. Two such steps exist today, the +multiprocessing step and the windowed reduce; both are pre-existing Python Arroyo +strategies wrapped unmodified rather than rewritten in Rust. + +The Python side does not implement an Arroyo strategy. It implements a deliberately simpler +interface, `RustOperatorDelegate`: `submit()` accepts and stores work, `poll()` processes +and **returns** results, `flush()` drains. Rust forwards whatever `poll` returned to the next +strategy, so the delegate never holds a reference to anything downstream. Delegates are +created by a `RustOperatorFactory` rather than passed in directly. + +## Alternatives + +- Reimplement multiprocessing and windowed reduce natively in Rust. +- Give the Python side a real Arroyo strategy interface with a handle to the next Rust + strategy. + +## Rationale + +This is meant to be a temporary solution to support more primitives quicker. +Ideally this will not survive when more native rust strategies are implemented. + +## Consequences + +- Returning instead of pushing makes the interface inherently asynchronous: a delegate may + be 1:1, 1:N, N:1 or N:0, which a per-input return signature cannot express. +- Exceptions out of `submit` are the delegate's control channel — `MessageRejected` becomes + backpressure, `InvalidMessage` routes to the DLQ, anything else panics. +- The factory is where state too expensive to rebuild lives — notably a pre-initialised + process pool — because Arroyo tears down and rebuilds the chain on every rebalance. +- Output messages are Arroyo "any messages", so a failure downstream of this operator cannot + be attributed to a specific offset. +- Payloads must pickle to reach a worker process, and the Rust `pyclass` types cannot, which + is why Python code sees the `Message` wrappers that build the Rust object lazily. + +## See also + +[The Python operator](../rust-arroyo/python-operator.md) · +[Rebalance and shutdown](../guarantees.md#rebalance-and-shutdown) diff --git a/sentry_streams/docs/architecture/decisions/0014-fuse-maps-into-segments.md b/sentry_streams/docs/architecture/decisions/0014-fuse-maps-into-segments.md new file mode 100644 index 00000000..e7fab441 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0014-fuse-maps-into-segments.md @@ -0,0 +1,49 @@ +# 0014 — Consecutive maps are fused into a chain; parallelism is a property of the segment + +**Status:** accepted · **Area:** Rust runtime + +## Decision + +The adapter does not translate steps one to one. Consecutive `Map` steps accumulate into a +chain, kept per route so each branch accumulates independently. Any non-map step closes the +open chain before it is added, keeping operator order faithful to pipeline order. A closed +chain becomes a single operator: a native `RuntimeOperator::Map` holding the fused function, +or — when the segment is configured with `multi_process` — a `PythonAdapter` wrapping +Arroyo's multiprocessing strategy, whose pool executes the fused function. + +Parallelism is configured only on a step that starts a segment; configuring it mid-chain is +rejected. The fused function is checked for picklability when the chain is finalised, **even +when multiprocessing is off**. + +## Alternatives + +- One operator per step, with a GIL acquisition and a message hop each. +- Check picklability only when a pool is actually configured. + +## Rationale + +This is meant to avoid going back and forth between Rust and python when the +back and forth is not needed. Moving messages around requires taking and releasing +the GIL, it is not the most efficient operation. + +It is also critical when we run steps across multiple processes as each parallel +step would need to have its own shared memory and processing pool. Chaining +stateless operations is a no brainer there. + +## Consequences + +- A run of maps costs one boundary crossing instead of one per step. +- A non-map step in the middle of what looks like one logical segment splits it: inserting a + filter between two maps produces two chains, two operators and potentially two pools. + Segmentation follows the pipeline's shape as much as the configuration. +- A pipeline cannot pass in a single-process deployment and then fail when someone enables a + pool in production; unpicklable functions are a build-time error naming the offending + steps. +- "Segment" means two different things — a config-list entry (one Kubernetes workload) and a + fusion boundary within a process. Both are in the [glossary](../README.md#glossary). + +## See also + +[Chaining and starts_segment](../rust-arroyo/adapter.md#chaining-and-starts_segment) · +[The adapter](../contracts.md#the-adapter) · +[The multiprocess step](../rust-arroyo/python-operator.md#the-multiprocess-step) diff --git a/sentry_streams/docs/architecture/decisions/0018-symmetric-buffered-metrics.md b/sentry_streams/docs/architecture/decisions/0018-symmetric-buffered-metrics.md new file mode 100644 index 00000000..f4a480b9 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/0018-symmetric-buffered-metrics.md @@ -0,0 +1,45 @@ +# 0018 — Metrics are symmetric across the language boundary and buffered, not sampled + +## Decision + +Step-level metrics are produced by an equivalent buffering component on each side of the +language boundary — `PipelineStats` in Python, thread-local state in Rust — with the same +semantics, the same metric names, the same `step` tag and the same ten-second flush +interval. The adapter wraps every application function it registers so entry, duration and +exceptions are recorded around the call. + +Aggregation is by buffering, not sampling. + +## Alternatives + +- Emit each metric as it happens. +- Sample in hot paths. +- Let each side report in whatever shape suits it. + +## Rationale + +Buffering is critical for metrics. Initially we were producing metrics at each +step, that ended up representing more than 90% of the CPU usage of the application. + +Moreover, as we process tens of thousands of messages per second per consumer +and we process each of them in multiple steps, even the data structure to +accumulate metrics is critical. + +PipelineStats is intentionally bare bone. Making it heavier can have disproportionate +effect. + +## Consequences + +- A step reports identically whether it was fused into a native operator or executed as a + Python function, so one set of dashboards reads a pipeline end to end. A new Rust strategy + is expected to keep that property. +- Rare events are kept rather than discarded, which sampling would not guarantee, and metrics + produced in tight loops do not cost a significant fraction of the work being measured. +- Metrics are delayed by up to the flush interval. +- Worker processes inherit nothing, so the pool initializer reconfigures metrics in each one + — which requires the config to be plain data rather than a live backend object + (an [invariant](../contracts.md#configuration)). + +## See also + +[Metrics](../rust-arroyo/metrics.md) diff --git a/sentry_streams/docs/architecture/decisions/README.md b/sentry_streams/docs/architecture/decisions/README.md new file mode 100644 index 00000000..37eb0ef2 --- /dev/null +++ b/sentry_streams/docs/architecture/decisions/README.md @@ -0,0 +1,43 @@ +# Decision log + +One page per architectural decision that shaped the platform. Each page states the +decision, what it rules out, and the consequences we live with. The **rationale** section +is where the reasoning behind the decision is recorded — why this option, at the time, over +the alternatives. + +These are records, not proposals. A decision that is later reversed keeps its page and gains +a superseded status, so the history of why the system looks like it does stays readable. + +The mechanics of each decision live in the architecture pages; a decision page links to them +rather than repeating them. + +## The platform shape + +| # | Decision | +| --- | --- | +| [0001](./0001-separate-application-from-deployment.md) | An application describes a dataflow; a separate config describes its deployment | +| [0002](./0002-adapters-as-the-runtime-boundary.md) | Runtimes are pluggable adapters with one method per primitive | +| [0004](./0004-complex-steps-as-sugar.md) | Higher level steps are sugar that converts to primitives, with a native override hook | +| [0005](./0005-statically-typed-dsl.md) | The DSL is statically typed end to end | + +## The Rust runtime + +| # | Decision | +| --- | --- | +| [0009](./0009-rust-runs-python-builds.md) | Rust owns the runtime loop; Python builds the pipeline | +| [0010](./0010-routes-for-branching.md) | Branching is expressed as routes on a single linear strategy chain | +| [0011](./0011-watermarks-drive-commits.md) | Commits are driven by periodic watermarks, not by data messages | +| [0013](./0013-embed-python-strategies-via-a-pull-delegate.md) | Existing Python Arroyo strategies are embedded through a pull-based delegate, not reimplemented | +| [0014](./0014-fuse-maps-into-segments.md) | Consecutive maps are fused into a chain; parallelism is a property of the segment | + +## Observability + +| # | Decision | +| --- | --- | +| [0018](./0018-symmetric-buffered-metrics.md) | Metrics are symmetric across the language boundary and buffered, not sampled | + +## Writing a new one + +Copy the shape of an existing page: a one-line status, **Decision**, **Alternatives**, +**Rationale**, **Consequences**, **See also**. Keep the decision to a few sentences and put +the reasoning in the rationale. Number sequentially; never renumber an existing page.