Cloudflare.Pipelines reference
LegacyPipeline
Section titled “LegacyPipeline”Source:
src/Cloudflare/Pipelines/LegacyPipeline.ts
A legacy Cloudflare Pipeline — the original HTTP-ingest → R2 batch
product (/accounts/{account}/pipelines).
A legacy pipeline accepts JSON events over HTTP (and/or a Worker
pipelines binding) and batches them into an R2 bucket using
S3-compatible credentials.
LegacyPipeline: Creating a Legacy Pipeline
Section titled “LegacyPipeline: Creating a Legacy Pipeline”HTTP ingest into R2
The S3-compatible credentials are derived from a Cloudflare API token: the access key id is the token id and the secret is the SHA-256 hex digest of the token value.
const bucket = yield* Cloudflare.R2.Bucket("events", {});
const pipeline = yield* Cloudflare.Pipelines.LegacyPipeline("ingest", { destination: { bucket: bucket.bucketName, credentials: { accessKeyId: yield* Config.Redacted("R2_ACCESS_KEY_ID"), secretAccessKey: yield* Config.Redacted("R2_SECRET_ACCESS_KEY"), }, },});// POST events to pipeline.endpointTuned batching and CORS
const pipeline = yield* Cloudflare.Pipelines.LegacyPipeline("ingest", { source: [ { type: "http", cors: { origins: ["https://example.com"] } }, ], destination: { bucket: bucket.bucketName, credentials, batch: { maxDurationS: 10, maxRows: 1000 }, compression: "gzip", prefix: "ingest", },});Pipeline
Section titled “Pipeline”Source:
src/Cloudflare/Pipelines/Pipeline.ts
A Cloudflare SQL Pipeline — the transform of the Pipelines product. A
pipeline is a single SQL statement that reads events from a
Stream and writes them to a Sink, both
referenced by name.
Also exported as Cloudflare.Basin.Pipeline.
The SQL is fixed at creation: changing it (or the name) triggers a
replacement. Nothing references a pipeline downstream, so replacements
are cheap. A changed statement is checked with Cloudflare’s SQL
validator during plan, so a syntax error fails the plan (with
PipelineSqlInvalid) instead of deleting the running pipeline.
References to streams/sinks that do not exist yet are not an error at
plan time — they are created in the same deploy.
Pipelines is ingest-only: to process data after the sink writes it,
subscribe to the sink’s bucket with Cloudflare.R2.BucketEventNotification.
Pipeline: Creating a Pipeline
Section titled “Pipeline: Creating a Pipeline”Stream → Sink passthrough
const stream = yield* Cloudflare.Basin.Stream("events", {});const sink = yield* Cloudflare.Basin.Sink("events-sink", { type: "r2", config: { bucket: bucket.bucketName, credentials },});
const pipeline = yield* Cloudflare.Basin.Pipeline("etl", { sql: Output.interpolate`INSERT INTO ${sink.name} SELECT * FROM ${stream.name}`,});Filtering transform
const pipeline = yield* Cloudflare.Basin.Pipeline("errors-only", { sql: Output.interpolate`INSERT INTO ${sink.name} SELECT * FROM ${stream.name} WHERE level = 'error'`,});Source:
src/Cloudflare/Pipelines/Sink.ts
A Cloudflare Pipelines sink — the destination of the Pipelines product
(also exported as Cloudflare.Basin.Sink). A SQL Pipeline
reads events from a Stream and writes them to a sink, which
stores them in R2 either as raw files (r2) or as an Iceberg table via
the Basin Catalog (basin_catalog, a.k.a. r2_data_catalog).
Sinks have no update API: changing the destination (bucket, path,
table, format, partitioning, rolling policy) triggers a replacement.
Equivalent spellings never replace — basin_catalog ⇄
r2_data_catalog, the config form ⇄ the catalog / table form,
omitted ⇄ explicit defaults. With engine-generated names a replacement
is seamless (the new sink gets a fresh name before the old one is
deleted); with an explicit name it collides, so prefer generated
names.
Credentials (credentials, token) are write-only and are not part of
the replacement diff, so rotating them never recreates the sink — a
catalog sink could not be recreated anyway, because Cloudflare refuses
to create sinks for existing Iceberg tables. The sink keeps the
credentials it was created with; to move it onto new ones, rename it
(which replaces it).
Pipelines is ingest-only — to react to files the sink writes, subscribe
to its bucket with Cloudflare.R2.BucketEventNotification.
Sink: Creating a Sink
Section titled “Sink: Creating a Sink”R2 sink with JSON output
The S3-compatible credentials are derived from a Cloudflare API token: the access key id is the token id and the secret is the SHA-256 hex digest of the token value.
const bucket = yield* Cloudflare.R2.Bucket("events", {});
const sink = yield* Cloudflare.Basin.Sink("events-sink", { type: "r2", config: { bucket, credentials: { accessKeyId: yield* Config.Redacted("R2_ACCESS_KEY_ID"), secretAccessKey: yield* Config.Redacted("R2_SECRET_ACCESS_KEY"), }, path: "ingest", rollingPolicy: { intervalSeconds: 30 }, },});Parquet output
const sink = yield* Cloudflare.Basin.Sink("parquet-sink", { type: "r2", config: { bucket, credentials }, format: { type: "parquet", compression: "zstd" },});Sink: Basin Catalog
Section titled “Sink: Basin Catalog”Iceberg table sink
const bucket = yield* Cloudflare.R2.Bucket("Lakehouse", {});const catalog = yield* Cloudflare.Basin.Catalog("Catalog", { bucket,});
const sink = yield* Cloudflare.Basin.Sink("PageViewTable", { type: "basin_catalog", catalog, table: { namespace: "web", name: "page_views" }, token: yield* Config.Redacted("CATALOG_TOKEN"),});Bucket-name form
const sink = yield* Cloudflare.Basin.Sink("iceberg-sink", { type: "r2_data_catalog", config: { bucket, tableName: "events", namespace: "default", token: yield* Config.Redacted("CATALOG_TOKEN"), },});Stream
Section titled “Stream”Source:
src/Cloudflare/Pipelines/Stream.ts
A Cloudflare Pipelines stream — the ingestion endpoint of the Pipelines
product (also exported as Cloudflare.Basin.Stream). Events are sent to
a stream from Workers (WriteStream / StreamSink) or over HTTP,
transformed by a SQL Pipeline, and written to a Sink.
The stream’s schema and format are fixed at creation (changing them
triggers a replacement, which drops buffered events and dependent
pipelines); the HTTP endpoint and Worker-binding toggles are mutable in
place.
Pipelines is ingest-only — there is no event source on a stream. To
process data after it lands, subscribe to the sink’s bucket with
Cloudflare.R2.BucketEventNotification.
Cloudflare accepts records that violate a structured stream’s schema (they are dropped later, during processing). Declare the schema as an Effect Schema to have the producer clients validate and encode every record before it is sent.
Stream: Creating a Stream
Section titled “Stream: Creating a Stream”Unstructured stream with default settings
const stream = yield* Cloudflare.Basin.Stream("events", {});Structured stream from an Effect Schema
class PageView extends Schema.Class<PageView>("PageView")({ url: Schema.String, at: Schema.Date, tags: Schema.Array(Schema.String), user: Schema.optional(Schema.Struct({ id: Schema.String })),}) {}
// Stream<PageView>: WriteStream / StreamSink are typed and encode recordsconst stream = yield* Cloudflare.Basin.Stream("PageViews", { schema: PageView,});Structured stream with an explicit field list
const stream = yield* Cloudflare.Basin.Stream("clicks", { schema: { fields: [ { type: "string", name: "url", required: true }, { type: "timestamp", name: "ts", unit: "millisecond" }, { type: "list", name: "tags", items: { type: "string" } }, { type: "struct", name: "user", fields: [{ type: "string", name: "id" }], }, ], },});Stream: HTTP ingestion
Section titled “Stream: HTTP ingestion”Authenticated endpoint
const stream = yield* Cloudflare.Basin.Stream("events", { http: true });// POST a JSON array to stream.endpoint with a `Pipelines Send` token,// or use Cloudflare.Basin.WriteStreamHttp / WriteStreamLocal.Public endpoint with CORS
const stream = yield* Cloudflare.Basin.Stream("beacons", { http: { enabled: true, authentication: false, cors: { origins: ["https://app.example.com"] }, },});Stream: Wiring into a Pipeline
Section titled “Stream: Wiring into a Pipeline”const pipeline = yield* Cloudflare.Basin.Pipeline("etl", { sql: Output.interpolate`INSERT INTO ${sink.name} SELECT * FROM ${stream.name}`,});StreamSink
Section titled “StreamSink”Source:
src/Cloudflare/Pipelines/StreamSink.ts
Binding service that exposes a Pipelines Stream as an Effect
Sink: run a Stream of records into it and every element is
ingested. Also exported as Cloudflare.Basin.StreamSink.
Records are typed, validated and encoded by the stream’s Effect Schema (plain JSON objects when it has none). Each upstream chunk is packed greedily into requests of at most 5 MB, preserving order.
Provide StreamSinkBinding (native Worker binding), StreamSinkHttp
(scoped Pipelines Send token) or StreamSinkLocal (current
credentials, for Actions) — each is the sink layered over the matching
WriteStream implementation.
StreamSink: Draining a Stream into a Pipelines Stream
Section titled “StreamSink: Draining a Stream into a Pipelines Stream”Run an Effect Stream into a Pipelines stream
export default Cloudflare.Worker( "Api", { main: import.meta.url }, Effect.gen(function* () { const sink = yield* Cloudflare.Basin.StreamSink(PageViews); return { fetch: Effect.gen(function* () { yield* EffectStream.fromIterable(views).pipe(EffectStream.run(sink)); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.Basin.StreamSinkBinding)),);Controlling request size
// at most 500 records per sendyield* events.pipe(EffectStream.rechunk(500), EffectStream.run(sink));StreamSinkBinding
Section titled “StreamSinkBinding”Source:
src/Cloudflare/Pipelines/StreamSinkBinding.tsKind: Layer · Provides:Cloudflare.Pipelines.StreamSink
Implementation of the StreamSink service over the native Worker pipelines binding (WriteStreamBinding). Registers the same binding as WriteStream, so a Worker can use both on one stream.
StreamSinkBinding: Providing the Layer
Section titled “StreamSinkBinding: Providing the Layer”Effect.gen(function* () { const sink = yield* Cloudflare.Basin.StreamSink(PageViews); // ...}).pipe(Effect.provide(Cloudflare.Basin.StreamSinkBinding));StreamSinkHttp
Section titled “StreamSinkHttp”Source:
src/Cloudflare/Pipelines/StreamSinkHttp.tsKind: Layer · Provides:Cloudflare.Pipelines.StreamSink
Implementation of the StreamSink service over the stream’s HTTP ingest endpoint (WriteStreamHttp), authenticated with a scoped Pipelines Send API token bound into the host. The stream needs http enabled.
StreamSinkHttp: Providing the Layer
Section titled “StreamSinkHttp: Providing the Layer”Effect.gen(function* () { const sink = yield* Cloudflare.Basin.StreamSink(PageViews); // ...}).pipe(Effect.provide(Cloudflare.Basin.StreamSinkHttp));StreamSinkLocal
Section titled “StreamSinkLocal”Source:
src/Cloudflare/Pipelines/StreamSinkLocal.tsKind: Layer · Provides:Cloudflare.Pipelines.StreamSink
Implementation of the StreamSink service over the stream’s HTTP ingest endpoint with the current credentials (WriteStreamLocal) — for draining an Effect Stream into a Pipelines stream from an Action, with no Worker host. The stream needs http enabled.
StreamSinkLocal: Providing the Layer
Section titled “StreamSinkLocal: Providing the Layer”const Seed = Alchemy.Action( "Seed", Effect.gen(function* () { const sink = yield* Cloudflare.Basin.StreamSink(PageViews); return Effect.fn(function* () { yield* EffectStream.fromIterable(views).pipe(EffectStream.run(sink)); }); }).pipe(Effect.provide(Cloudflare.Basin.StreamSinkLocal)),);WriteStream
Section titled “WriteStream”Source:
src/Cloudflare/Pipelines/WriteStream.ts
Binding service that turns a Pipelines Stream (or a
LegacyPipeline) into a WriteStreamClient you can call
from a Worker’s (or Action’s) runtime Effect. Also exported as
Cloudflare.Basin.WriteStream.
Producer-only — send ingests a batch of records into the stream. When
the stream was declared with an Effect Schema, the client is typed by
it (send(records: PageView[])) and every record is validated and
encoded with Schema.toCodecJson before it is sent — Cloudflare
accepts schema-violating records and silently drops them later, so
this is the only validation. Without a schema, send takes plain JSON
objects.
Provide one implementation layer:
WriteStreamBinding— native WorkerpipelinesbindingWriteStreamHttp— the stream’s HTTP ingest endpoint with a scopedPipelines SendAPI token (the stream needshttpenabled)WriteStreamLocal— the HTTP ingest endpoint with the current credentials, for Actions and other deploy-time Effects
WriteStream: Sending Events
Section titled “WriteStream: Sending Events”Typed producer route
class PageView extends Schema.Class<PageView>("PageView")({ url: Schema.String, at: Schema.Date,}) {}
export const PageViews = Cloudflare.Basin.Stream("PageViews", { schema: PageView,});
export default Cloudflare.Worker( "Api", { main: import.meta.url }, Effect.gen(function* () { const views = yield* Cloudflare.Basin.WriteStream(PageViews); return { fetch: Effect.gen(function* () { yield* views.send([new PageView({ url: "/", at: new Date() })]); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.Basin.WriteStreamBinding)),);Untyped stream
const events = yield* Cloudflare.Basin.WriteStream(Events);yield* events.send([{ event: "click", at: new Date().toISOString() }]);WriteStream: Handling Errors
Section titled “WriteStream: Handling Errors”yield* views.send(batch).pipe( Effect.catchTag("StreamSendError", (e) => Effect.logWarning(`dropped ${batch.length} events: ${e.message}`), ),);WriteStreamBinding
Section titled “WriteStreamBinding”Source:
src/Cloudflare/Pipelines/WriteStreamBinding.tsKind: Layer · Provides:Cloudflare.Pipelines.WriteStream
Implementation of the WriteStream service over a native Worker
pipelines binding. Registers the binding on the host Worker and, when
the stream has an Effect Schema, encodes records with it before
send.
WriteStreamBinding: Providing the Layer
Section titled “WriteStreamBinding: Providing the Layer”Effect.gen(function* () { const views = yield* Cloudflare.Basin.WriteStream(PageViews); // ...}).pipe(Effect.provide(Cloudflare.Basin.WriteStreamBinding));WriteStreamHttp
Section titled “WriteStreamHttp”Source:
src/Cloudflare/Pipelines/WriteStreamHttp.tsKind: Layer · Provides:Cloudflare.Pipelines.WriteStream
HTTP-backed implementation of the WriteStream service. Creates
a scoped AccountApiToken with the Pipelines Send permission, binds
it into the host, and POSTs records to the stream’s HTTP ingest
endpoint (https://{stream_id}.ingest.cloudflare.com).
The stream must have its HTTP endpoint enabled — declare it with
http: true (authenticated). Records are validated and encoded with
the stream’s Effect Schema exactly as with WriteStreamBinding.
WriteStreamHttp: Providing the Layer
Section titled “WriteStreamHttp: Providing the Layer”const Events = Cloudflare.Basin.Stream("Events", { http: true });
Effect.gen(function* () { const events = yield* Cloudflare.Basin.WriteStream(Events); // ...}).pipe(Effect.provide(Cloudflare.Basin.WriteStreamHttp));WriteStreamLocal
Section titled “WriteStreamLocal”Source:
src/Cloudflare/Pipelines/WriteStreamLocal.tsKind: Layer · Provides:Cloudflare.Pipelines.WriteStream
Local implementation of the WriteStream service — POSTs records
to the stream’s HTTP ingest endpoint with the current credentials
instead of a native Worker binding (WriteStreamBinding) or a scoped
API token (WriteStreamHttp).
Provide it on an Action (or any deploy-time Effect) to send events with
the same typed client you use inside a Worker. The stream needs its HTTP
endpoint enabled (http: true), and the current credentials need the
Pipelines Send permission when it is authenticated.
WriteStreamLocal: Providing the Layer
Section titled “WriteStreamLocal: Providing the Layer”const Seed = Alchemy.Action( "Seed", Effect.gen(function* () { const views = yield* Cloudflare.Basin.WriteStream(PageViews); return Effect.fn(function* () { yield* views.send([new PageView({ url: "/", at: new Date() })]); }); }).pipe(Effect.provide(Cloudflare.Basin.WriteStreamLocal)),);