Cloudflare.K2 reference
ReadSubscription
Section titled “ReadSubscription”Source:
src/Cloudflare/K2/ReadSubscription.ts
Binding service that turns a K2 Subscription into a pull client:
lease a batch of records, then acknowledge it (advance the subscription),
nack it (redeliver), or extend its five-minute lease.
The native k2 Worker binding can only produce, so this service ships
HTTP implementations only: ReadSubscriptionHttp (scoped
K2 Consume token) and ReadSubscriptionLocal (current
credentials). For a long-running consumer loop, see
consumeStreamRecords.
ReadSubscription: Pulling Records
Section titled “ReadSubscription: Pulling Records”const inbox = yield* Cloudflare.K2.ReadSubscription(Analytics);
const batch = yield* inbox.consume({ workerId: "worker-1", maxRecords: 100 });if (Option.isSome(batch)) { yield* Effect.forEach(batch.value.records, (record) => Effect.log(new TextDecoder().decode(record.content)), ); yield* inbox.ack(batch.value);}ReadSubscription: Redelivery
Section titled “ReadSubscription: Redelivery”yield* process(batch.records).pipe( Effect.andThen(inbox.ack(batch)), Effect.catch(() => inbox.nack(batch)),);ReadSubscriptionHttp
Section titled “ReadSubscriptionHttp”Source:
src/Cloudflare/K2/ReadSubscriptionHttp.tsKind: Layer · Provides:Cloudflare.K2.ReadSubscription
HTTP-backed implementation of the ReadSubscription service. Mints
a scoped K2 Consume AccountApiToken, binds it into the host, and
leases batches through the stream’s data-plane API
(https://<id>.k2.cloudflarestorage.com/subscriptions/...).
ReadSubscriptionHttp: Providing the Layer
Section titled “ReadSubscriptionHttp: Providing the Layer”Effect.gen(function* () { const inbox = yield* Cloudflare.K2.ReadSubscription(Analytics); // ...}).pipe(Effect.provide(Cloudflare.K2.ReadSubscriptionHttp));ReadSubscriptionLocal
Section titled “ReadSubscriptionLocal”Source:
src/Cloudflare/K2/ReadSubscriptionLocal.tsKind: Layer · Provides:Cloudflare.K2.ReadSubscription
Local implementation of the ReadSubscription service — leases
batches over the K2 subscription HTTP API with the current
credentials instead of a scoped token. Provide it on an Action (or any
deploy-time Effect), or on a polling event source run in-process.
The subscription’s ids are resolved at apply time, so the subscription can be created in the same deploy.
ReadSubscriptionLocal: Providing the Layer
Section titled “ReadSubscriptionLocal: Providing the Layer”const Drain = Alchemy.Action( "Drain", Effect.gen(function* () { const inbox = yield* Cloudflare.K2.ReadSubscription(Analytics); return Effect.fn(function* () { const batch = yield* inbox.consume({ workerId: "drain" }); if (Option.isSome(batch)) yield* inbox.ack(batch.value); }); }).pipe(Effect.provide(Cloudflare.K2.ReadSubscriptionLocal)),);Stream
Section titled “Stream”Source:
src/Cloudflare/K2/Stream.ts
A Cloudflare K2 stream — a durable, ordered log of records. Producers
append records over HTTP or from Workers through a k2 binding;
consumers read them through a Subscription, each with its own
position in the stream.
The name is fixed at creation (changing it replaces the stream); retention and the HTTP / Worker-binding inputs are mutable in place.
Stream: Creating a Stream
Section titled “Stream: Creating a Stream”Stream with default settings
// Worker-binding input enabled, HTTP input disabled, 7-day retention.const orders = yield* Cloudflare.K2.Stream("Orders");Custom retention
const orders = yield* Cloudflare.K2.Stream("Orders", { retention: "30 days",});Stream: Typed Records
Section titled “Stream: Typed Records”const Order = Schema.Struct({ orderId: Schema.String, amountCents: Schema.Number });
// Producers send `Order`s; consumers receive them decoded in `record.value`.const orders = yield* Cloudflare.K2.Stream("Orders", { schema: Order });Stream: HTTP Input
Section titled “Stream: HTTP Input”Authenticated HTTP input
// POST records to orders.endpoint + "/produce" with a K2 Produce token.const orders = yield* Cloudflare.K2.Stream("Orders", { http: true,});Public HTTP input for browsers
const clicks = yield* Cloudflare.K2.Stream("Clicks", { http: { authentication: false, cors: ["https://app.example.com"] }, workerBinding: false,});Stream: Producing from a Worker
Section titled “Stream: Producing from a Worker”export default Cloudflare.Worker( "Api", { main: import.meta.url }, Effect.gen(function* () { const orders = yield* Cloudflare.K2.WriteStream(Orders); return { fetch: Effect.gen(function* () { yield* orders.send([{ content: JSON.stringify({ id: 1 }) }]); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.K2.WriteStreamBinding)),);Stream: Consuming
Section titled “Stream: Consuming”const analytics = yield* Cloudflare.K2.Subscription("Analytics", { stream: orders, startAt: "earliest",});StreamEventSourcePolling
Section titled “StreamEventSourcePolling”Source:
src/Cloudflare/K2/StreamEventSourcePoller.tsKind: Layer · Provides:Cloudflare.K2.StreamEventSource
Polling implementation of StreamEventSource for hosts with a
long-lived process (any host providing ServerHost: Fly, ECS, EC2,
Railway, Hetzner, Docker). Requires a ReadSubscription implementation.
StreamEventSourcePolling: Providing the Layer
Section titled “StreamEventSourcePolling: Providing the Layer”Effect.gen(function* () { yield* Cloudflare.K2.consumeStreamRecords(Orders, handler);}).pipe( Effect.provide( Layer.provideMerge( Cloudflare.K2.StreamEventSourcePolling, Cloudflare.K2.ReadSubscriptionHttp, ), ),);StreamSink
Section titled “StreamSink”Source:
src/Cloudflare/K2/StreamSink.ts
Binding service that exposes a K2 Stream as an Effect Sink: run
a Stream of records into it and every element is appended.
Each upstream chunk is packed, in order, into send calls of at most
5 MB; a record over the 1 MB limit ships alone so K2’s rejection
surfaces. A batch that failed with K2Unavailable (not stored) is
retried a bounded number of times; K2AppendOutcomeUnknown is never
retried because the batch may have been stored.
Provide StreamSinkBinding (native Worker binding),
StreamSinkHttp (scoped K2 Produce token) or
StreamSinkLocal (current credentials, for Actions).
StreamSink: Draining a Stream into K2
Section titled “StreamSink: Draining a Stream into K2”Run a Stream into a K2 stream
export default Cloudflare.Worker( "Api", { main: import.meta.url }, Effect.gen(function* () { const sink = yield* Cloudflare.K2.StreamSink(Orders); return { fetch: Effect.gen(function* () { yield* Stream.fromIterable(events).pipe( Stream.map((event) => ({ content: JSON.stringify(event) })), Stream.run(sink), ); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.K2.StreamSinkBinding)),);Typed sink with a Schema
const sink = yield* Cloudflare.K2.StreamSink(Orders, { schema: Order });yield* Stream.fromIterable(orders).pipe(Stream.run(sink));StreamSinkBinding
Section titled “StreamSinkBinding”Source:
src/Cloudflare/K2/StreamSinkBinding.tsKind: Layer · Provides:Cloudflare.K2.StreamSink
Implementation of the StreamSink service over the native Worker
k2 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”export default Cloudflare.Worker( "Api", { main: import.meta.url }, Effect.gen(function* () { const sink = yield* Cloudflare.K2.StreamSink(Orders); return { fetch: Effect.gen(function* () { yield* Stream.make({ content: "a" }, { content: "b" }).pipe(Stream.run(sink)); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.K2.StreamSinkBinding)),);StreamSinkHttp
Section titled “StreamSinkHttp”Source:
src/Cloudflare/K2/StreamSinkHttp.tsKind: Layer · Provides:Cloudflare.K2.StreamSink
Implementation of the StreamSink service over the stream’s HTTP
input (WriteStreamHttp), authenticated with a scoped K2 Produce
API token bound into the host. The stream needs http: true.
StreamSinkHttp: Providing the Layer
Section titled “StreamSinkHttp: Providing the Layer”Effect.gen(function* () { const sink = yield* Cloudflare.K2.StreamSink(Orders); // ...}).pipe(Effect.provide(Cloudflare.K2.StreamSinkHttp));StreamSinkLocal
Section titled “StreamSinkLocal”Source:
src/Cloudflare/K2/StreamSinkLocal.tsKind: Layer · Provides:Cloudflare.K2.StreamSink
Implementation of the StreamSink service that appends over the
stream’s HTTP input with the current credentials
(WriteStreamLocal) — for draining a Stream into K2 from an Action
or other deploy-time Effect. The stream needs http: true.
StreamSinkLocal: Providing the Layer
Section titled “StreamSinkLocal: Providing the Layer”const Seed = Alchemy.Action( "Seed", Effect.gen(function* () { const sink = yield* Cloudflare.K2.StreamSink(Orders); return Effect.fn(function* () { yield* Stream.range(1, 500).pipe( Stream.map((id) => ({ content: JSON.stringify({ id }) })), Stream.run(sink), ); }); }).pipe(Effect.provide(Cloudflare.K2.StreamSinkLocal)),);Subscription
Section titled “Subscription”Source:
src/Cloudflare/K2/Subscription.ts
A K2 subscription — a named read position on a Stream. Consumers
lease batches of records from a subscription, then acknowledge them to
advance it. Every subscription on a stream reads every record
independently; consumers sharing one subscription compete for batches.
Subscriptions are immutable: changing the stream, name, or start position replaces the subscription.
Subscription: Creating a Subscription
Section titled “Subscription: Creating a Subscription”Read new records
const orders = yield* Cloudflare.K2.Stream("Orders");const analytics = yield* Cloudflare.K2.Subscription("Analytics", { stream: orders,});Read from the oldest retained record
const backfill = yield* Cloudflare.K2.Subscription("Backfill", { stream: orders, startAt: "earliest",});Subscription: Pulling Records
Section titled “Subscription: Pulling Records”const inbox = yield* Cloudflare.K2.ReadSubscription(analytics);const batch = yield* inbox.consume({ workerId: "worker-1", maxRecords: 100 });if (Option.isSome(batch)) { yield* Effect.forEach(batch.value.records, handle); yield* inbox.ack(batch.value);}WriteStream
Section titled “WriteStream”Source:
src/Cloudflare/K2/WriteStream.ts
Binding service that turns a K2 Stream into a typed
WriteStreamClient you can call from a Worker’s (or any host’s)
runtime Effect.
send appends a batch of records atomically: all of them are stored, or
none. A batch can be at most 5 MB and each record at most 1 MB.
WriteStream: Sending Records
Section titled “WriteStream: Sending Records”const orders = yield* Cloudflare.K2.WriteStream(Orders);
return { fetch: Effect.gen(function* () { yield* orders.send([ { content: JSON.stringify({ id: 1 }), headers: { source: "api" } }, { content: new Uint8Array([1, 2, 3]) }, ]); return HttpServerResponse.empty({ status: 202 }); }),};WriteStream: Typed Records
Section titled “WriteStream: Typed Records”const Order = Schema.Struct({ id: Schema.Number, total: Schema.Number });const Orders = Cloudflare.K2.Stream("Orders", { schema: Order });
const orders = yield* Cloudflare.K2.WriteStream(Orders);// JSON-encoded, sent with `content-type: application/json`yield* orders.send([{ id: 1, total: 42 }]);WriteStream: Handling Errors
Section titled “WriteStream: Handling Errors”yield* orders.send(records).pipe( // K2Unavailable: the batch was NOT stored. Never retry // K2AppendOutcomeUnknown — the batch may have been stored. Effect.retry({ while: (e) => e._tag === "K2Unavailable", times: 3, }),);Provide WriteStreamBinding (native k2 Worker binding),
WriteStreamHttp (scoped K2 Produce token over the stream’s HTTP
input) or WriteStreamLocal (current credentials, for Actions).
WriteStreamBinding
Section titled “WriteStreamBinding”Source:
src/Cloudflare/K2/WriteStreamBinding.tsKind: Layer · Provides:Cloudflare.K2.WriteStream
Implementation of the WriteStream service over the native Worker
k2 binding. The stream must have its Worker-binding input enabled (the
default).
K2 has no local simulation, so a Worker using this layer cannot run under
alchemy dev.
WriteStreamBinding: Providing the Layer
Section titled “WriteStreamBinding: Providing the Layer”export default Cloudflare.Worker( "Api", { main: import.meta.url }, Effect.gen(function* () { const orders = yield* Cloudflare.K2.WriteStream(Orders); return { fetch: Effect.gen(function* () { yield* orders.send([{ content: "hello" }]); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.K2.WriteStreamBinding)),);WriteStreamHttp
Section titled “WriteStreamHttp”Source:
src/Cloudflare/K2/WriteStreamHttp.tsKind: Layer · Provides:Cloudflare.K2.WriteStream
HTTP-backed implementation of the WriteStream service. Mints a
scoped K2 Produce AccountApiToken, binds it into the host, and
appends records through the stream’s HTTP input
(POST https://<id>.k2.cloudflarestorage.com/produce).
The stream must have its HTTP input enabled (http: true). Works on any
host — Workers, containers, Lambda.
WriteStreamHttp: Providing the Layer
Section titled “WriteStreamHttp: Providing the Layer”const Orders = Cloudflare.K2.Stream("Orders", { http: true });
Effect.gen(function* () { const orders = yield* Cloudflare.K2.WriteStream(Orders); // ...}).pipe(Effect.provide(Cloudflare.K2.WriteStreamHttp));WriteStreamLocal
Section titled “WriteStreamLocal”Source:
src/Cloudflare/K2/WriteStreamLocal.tsKind: Layer · Provides:Cloudflare.K2.WriteStream
Local implementation of the WriteStream service — appends records
through the stream’s HTTP input with the current credentials instead
of a native binding or a scoped token. Provide it on an Action (or any
deploy-time Effect) to seed a stream with the same client a Worker uses.
The stream must have its HTTP input enabled (http: true). The stream id
is resolved at apply time, so the stream can be created in the same
deploy.
WriteStreamLocal: Providing the Layer
Section titled “WriteStreamLocal: Providing the Layer”const Seed = Alchemy.Action( "Seed", Effect.gen(function* () { const orders = yield* Cloudflare.K2.WriteStream(Orders); return Effect.fn(function* () { yield* orders.send([{ content: "hello" }]); }); }).pipe(Effect.provide(Cloudflare.K2.WriteStreamLocal)),);