Skip to content

Cloudflare.K2 reference

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.

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);
}
yield* process(batch.records).pipe(
Effect.andThen(inbox.ack(batch)),
Effect.catch(() => inbox.nack(batch)),
);

Source: src/Cloudflare/K2/ReadSubscriptionHttp.ts Kind: 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/...).

Effect.gen(function* () {
const inbox = yield* Cloudflare.K2.ReadSubscription(Analytics);
// ...
}).pipe(Effect.provide(Cloudflare.K2.ReadSubscriptionHttp));

Source: src/Cloudflare/K2/ReadSubscriptionLocal.ts Kind: 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)),
);

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 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",
});
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 });

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,
});
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)),
);
const analytics = yield* Cloudflare.K2.Subscription("Analytics", {
stream: orders,
startAt: "earliest",
});

Source: src/Cloudflare/K2/StreamEventSourcePoller.ts Kind: 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,
),
),
);

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).

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));

Source: src/Cloudflare/K2/StreamSinkBinding.ts Kind: 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.

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)),
);

Source: src/Cloudflare/K2/StreamSinkHttp.ts Kind: 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.

Effect.gen(function* () {
const sink = yield* Cloudflare.K2.StreamSink(Orders);
// ...
}).pipe(Effect.provide(Cloudflare.K2.StreamSinkHttp));

Source: src/Cloudflare/K2/StreamSinkLocal.ts Kind: 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.

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)),
);

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.

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",
});
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);
}

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.

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 });
}),
};
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 }]);
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).

Source: src/Cloudflare/K2/WriteStreamBinding.ts Kind: 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.

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)),
);

Source: src/Cloudflare/K2/WriteStreamHttp.ts Kind: 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.

const Orders = Cloudflare.K2.Stream("Orders", { http: true });
Effect.gen(function* () {
const orders = yield* Cloudflare.K2.WriteStream(Orders);
// ...
}).pipe(Effect.provide(Cloudflare.K2.WriteStreamHttp));

Source: src/Cloudflare/K2/WriteStreamLocal.ts Kind: 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.

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)),
);