The record shape this pipeline carries.
Constructs a PipelineBinding instance.
The record shape this pipeline carries.
Optionalsettings: PipelineSettings
Resolved settings. Defaults to resolve('pipeline'), so callers
normally construct this with no arguments at all.
The binding name this instance resolves.
The binding name as declared in wrangler.json.
Checks whether the binding is available on the given environment.
Use this to degrade gracefully when a binding is optional.
The Worker environment.
true when the binding is present.
Resolves the binding, failing fast when it is not configured.
The Worker environment.
The resolved binding.
MissingBindingError When the binding is absent from the environment.
Sends records to the pipeline.
Resolves when the pipeline has accepted the records, not when they have been written — the batching that makes a pipeline worth using happens after this returns, so a successful send is a promise to land the records rather than evidence that it happened.
An empty array is sent as given rather than short-circuited, because whether an empty batch is worth a call is the pipeline's business and a caller that filtered everything out has already decided to send.
MissingBindingError When the pipeline binding is absent.
Sends records without letting a failure reach the caller.
true when the pipeline accepted the records.
The method an archival sink should call. A missing binding, a rejected batch, an unreachable pipeline — each is logged and reported through the return value, so the work that produced the records finishes either way.
Reach for send instead when the pipeline is the only copy: silence is the right answer for a second sink and the wrong one for the only sink.
Accessor for a Pipelines binding.
Remarks
A pipeline takes structured records and lands them somewhere durable — R2, as files, batched and partitioned by rules configured on the pipeline rather than here. The Worker sends records and stops thinking about them.
That last part is the whole point, and it is what separates this from R2Binding. Writing to R2 directly means a producer owns buffering: how many records make a file, how long to wait for a slow hour, what to do with a partial batch when the isolate is evicted. Ingestion that arrives in bursts — a large batch produced inside one cron window — are exactly where that goes wrong, and getting it wrong loses data quietly.
It is also not QueueBinding, though both take a batch and return quickly. A queue exists so another Worker can act on each message; a pipeline exists so nothing has to. If the records are going to be read as files later rather than processed one at a time, the queue's consumer is a step that only exists to write them down.
Two ways to send, and the safe one is the point
send throws, like every other accessor here. sendSafe logs and returns
false, mirroring AnalyticsBinding'swriteSafe— and for the same reason. A pipeline's usual job is a second, independent sink beside the one doing the real work: archiving the raw record a provider returned while the coerced one goes on to storage. A sink like that failing the run that produced the records inverts its purpose, because the archive exists to explain a bad run and would instead be causing one.Which to call follows from whether the record has anywhere else to be. An audit trail that is the only copy should throw.
Why there is no
PipelineService, and no transformation supportNo kernel capability to implement, as with WorkflowBinding — and a pipeline is closer to plumbing than to a capability in any case.
Transformation is a separate matter. A pipeline can run records through a
PipelineTransformationEntrypointbefore landing them, which — like a Workflow'sWorkflowEntrypoint— is a class the consuming Worker exports and names in its ownwrangler.json. A base class that must be exported from someone else's entry has nothing an adapter can add, so this package wraps the producer side only.Example
Author
Bayu Dwiyan Satria
Version
1.0.0
Since
1.3.0