Skip to content

feat(cloudflare)!: K2 streams and Basin (Pipelines, Catalog, Tables) - #2010

Merged
sam-goodwin merged 12 commits into
mainfrom
feat/cloudflare-k2-basin
Oct 6, 2026
Merged

sam-goodwin merged 12 commits into
mainfrom
feat/cloudflare-k2-basin

Conversation

@sam-goodwin

@sam-goodwin sam-goodwin commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Cloudflare K2 streams and a Cloudflare.Basin namespace (Pipelines + Catalog + Iceberg tables), with hand-written guides for both. Builds on alchemy-run/distilled#730 (merged), which main's distilled pin already includes.

Breaking: a Pipelines.Stream declared without http no longer gets a public, unauthenticated ingest endpoint. The next deploy closes it in place, without replacing the stream. To keep it public, set http: { enabled: true, authentication: false }.

K2: a durable log you consume yourself

The schema lives on the stream, so every binding on it is typed the same way:

const Order = Schema.Struct({ orderId: Schema.String, amountCents: Schema.Number });

export const Orders = Cloudflare.K2.Stream("Orders", {
  schema: Order,                // Stream<Order>; a client-side codec, never persisted
  retention: "7 days",          // updated in place
  // http: true,                // HTTP /produce input, off by default; `true` = authenticated
});

Produce from an Effect Worker:

export default class Api extends Cloudflare.Worker<Api>()(
  "Api",
  { main: import.meta.url },
  Effect.gen(function* () {
    const orders = yield* Cloudflare.K2.WriteStream(Orders);   // send(records: Order[])
    const sink = yield* Cloudflare.K2.StreamSink(Orders);      // Sink<void, Order>, packs ≤ 5 MB requests

    return {
      fetch: Effect.gen(function* () {
        const order = yield* HttpServerRequest.schemaBodyJson(Order);
        yield* orders.send([order]).pipe(
          Effect.catchTag("K2Unavailable", () => Effect.fail(new Retry())), // not stored, safe to retry
          // K2AppendOutcomeUnknown: may have been stored, never retried automatically
        );
        return HttpServerResponse.empty({ status: 202 });
      }),
    };
  }).pipe(Effect.provide(Layer.mergeAll(Cloudflare.K2.WriteStreamBinding, Cloudflare.K2.StreamSinkBinding))),
) {}

Swap the layer to produce from somewhere else:

Effect.provide(Cloudflare.K2.WriteStreamHttp)   // Lambda / container: mints a scoped "K2 Produce" token
Effect.provide(Cloudflare.K2.WriteStreamLocal)  // Actions / scripts: current credentials

Consume from a long-running host (Fly, ECS, EC2, Railway, Hetzner, Docker). Records arrive decoded with the stream's schema. Success acks the batch, failure nacks it for redelivery, and the lease is extended while the handler runs:

yield* Cloudflare.K2.consumeStreamRecords(
  Orders,
  { startAt: "earliest", maxRecords: 500, concurrency: 4 },
  (records) =>
    records.pipe(
      Stream.map((record) => record.value),   // Order
      Stream.filter((order) => order.amountCents > 100_000),
      Stream.runForEach(flagForReview),
    ),
).pipe(
  Effect.provide(
    Layer.provideMerge(Cloudflare.K2.StreamEventSourcePolling, Cloudflare.K2.ReadSubscriptionHttp),
  ),
);

Each consuming host gets its own subscription (fan-out). To split records between hosts, pass a shared Subscription:

export const Fraud = Cloudflare.K2.Subscription("Fraud", { stream: Orders, startAt: "latest" });

yield* Cloudflare.K2.consumeStreamRecords(Orders, { subscription: Fraud }, handler);

Or drive the leases yourself:

const inbox = yield* Cloudflare.K2.ReadSubscription(Fraud);
const batch = yield* inbox.consume({ workerId: "w1", maxRecords: 100 });   // Option<Batch>
if (Option.isSome(batch)) {
  yield* process(batch.value.records).pipe(
    Effect.andThen(inbox.ack(batch.value)),
    Effect.catchCause(() => inbox.nack(batch.value)),
  );
}

Async Workers get a typed binding:

yield* Cloudflare.Worker("AsyncApi", { main: "./src/async.ts", env: { ORDERS: Orders } });

// src/async.ts — env.ORDERS: Cloudflare.K2.K2StreamBinding
const result = await env.ORDERS.send([{ content: new TextEncoder().encode(JSON.stringify(order)) }]);
if (!result.success) return new Response(result.error.message, { status: result.error.retryable ? 503 : 500 });

Workers can't consume K2 yet, because K2 can't push to them; push consumers are on Cloudflare's roadmap. No Worker consumer layer ships, and the StreamEventSource docs have a TODO for it.

Basin: Pipelines + Catalog + Iceberg tables

Cloudflare.Basin.* holds the same values as Cloudflare.Pipelines.* and Cloudflare.R2.DataCatalog. The old paths and persisted resource types are unchanged.

const PageView = Schema.Struct({
  userId: Schema.String,
  path: Schema.String,
  at: Schema.Date,
  tags: Schema.optional(Schema.Array(Schema.String)),   // list fields are new
});

const lake    = yield* Cloudflare.R2.Bucket("Lake");
const catalog = yield* Cloudflare.Basin.Catalog("Catalog", { bucket: lake });

// The Effect Schema becomes Cloudflare's field list; an equivalent { fields } is a no-op
const views = yield* Cloudflare.Basin.Stream("PageViews", { schema: PageView });

// The catalog resource is passed in, so the sink deploys after the catalog is enabled
const table = yield* Cloudflare.Basin.Sink("PageViewTable", {
  type: "basin_catalog",
  catalog,
  table: { namespace: "web", name: "page_views" },
  token,                                                // rotating it no longer replaces the sink
});

// SQL is validated during plan, so bad SQL fails before anything is replaced
yield* Cloudflare.Basin.Pipeline("PageViewEtl", {
  sql: Output.interpolate`INSERT INTO ${table.name} SELECT * FROM ${views.name}`,
});

Producers work like K2's and are typed by the stream's schema:

const writer = yield* Cloudflare.Basin.WriteStream(PageViews);   // send(records: PageView[])
const sink   = yield* Cloudflare.Basin.StreamSink(PageViews);    // Sink<void, PageView>
// .pipe(Effect.provide(Cloudflare.Basin.WriteStreamBinding))  |  WriteStreamHttp ("Pipelines Send")  |  WriteStreamLocal

Tables you own directly go over Iceberg REST, using @distilled.cloud/iceberg:

const orders = yield* Cloudflare.Basin.Table("Orders", {
  catalog,
  namespace: "sales",
  schema: Order,                     // additive changes evolve in place; destructive changes fail the plan
  partitionBy: ["day(at)"],
  properties: { "write.format.default": "parquet" },
  maintenance: { compaction: { targetSizeMb: "256" } },
  delete: "retain",                  // default; also "drop" | "purge"
});
// orders.identifier → "sales.orders", orders.tableUuid, orders.metadataLocation

Props that reference another resource take the resource or its identifier (bucket: string | R2.Bucket, stream: string | K2.Stream, catalog: string | Basin.Catalog). R2.DataCatalog's bucketName still works as a deprecated alias.

Docs

  • New guide: Cloudflare → Messaging & events → K2 streams (/cloudflare/messaging/k2).
  • New guide: Cloudflare → Data → Pipelines (/cloudflare/data/pipelines): stream → SQL → Iceberg or R2-file sink.
  • New guide: Cloudflare → Data → Iceberg tables (/cloudflare/data/iceberg-tables): the catalog, Basin.Table, and querying from other engines.
  • The Cloudflare overview links all three.

Dependencies

  • Adds @distilled.cloud/iceberg.
  • The Prisma example pins typescript, so the lockfile keeps resolving against the workspace TypeScript.

@alchemy-version-bot

alchemy-version-bot Bot commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Install the packages built from this commit:

alchemy

pnpm install https://pkg.alchemy.run/alchemy/pr:2010:f5e2556
@alchemy.run (6)
pnpm install https://pkg.alchemy.run/@alchemy.run/better-auth/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@alchemy.run/cloudflare-runtime/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@alchemy.run/frontend-frameworks/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@alchemy.run/node-utils/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@alchemy.run/floci/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@alchemy.run/pkg/pr:2010:f5e2556
@distilled.cloud (16)
pnpm install https://pkg.alchemy.run/@distilled.cloud/core/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/acme/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/aws/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/axiom/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/cloudflare/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/doppler/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/fly-io/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/github/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/hetzner/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/infisical/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/neon/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/prisma/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/planetscale/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/railway/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/stripe/pr:2010:f5e2556
pnpm install https://pkg.alchemy.run/@distilled.cloud/zerossl/pr:2010:f5e2556

Published Oct 6, 2026, 11:04 PM UTC. Expires Oct 13, 2026, 11:04 PM UTC, extended while this pull request is open.

@sam-goodwin
sam-goodwin force-pushed the feat/cloudflare-k2-basin branch 3 times, most recently from b3bae4b to e6ebcce Compare October 6, 2026 04:37
@sam-goodwin
sam-goodwin marked this pull request as ready for review October 6, 2026 08:33
@sam-goodwin
sam-goodwin force-pushed the feat/cloudflare-k2-basin branch 2 times, most recently from 250d4f6 to b0d2fc6 Compare October 6, 2026 11:04
The Prisma example pins `typescript` so @prisma/orm-postgres keeps
resolving against the workspace TypeScript after the lockfile refresh.
- Cloudflare.K2: Stream, Subscription, WriteStream (Binding/Http/Local),
  StreamSink, ReadSubscription, consumeStreamRecords + StreamEventSourcePolling
- Cloudflare.Basin: Pipelines + Catalog re-exports, Basin.Table over Iceberg REST
- Pipelines: typed WriteStream, Http/Local producers, StreamSink, Effect Schema
  and list/struct stream schemas, basin_catalog sinks taking the catalog,
  plan-time SQL validation

BREAKING: a Pipelines Stream declared without `http` no longer exposes a
public ingest endpoint; set `http: { enabled: true, authentication: false }`
to keep it public.
@sam-goodwin
sam-goodwin force-pushed the feat/cloudflare-k2-basin branch from b0d2fc6 to 44a6d9b Compare October 6, 2026 16:24
…ePoller.ts

StreamEventSource.ts keeps the contract and consumeStreamRecords; the
StreamEventSourcePolling layer and its pull loop (formerly StreamPoller.ts)
live together in StreamEventSourcePoller.ts.
K2.Stream("Orders", { schema: Order }) returns a Stream<Order>; WriteStream,
StreamSink and consumeStreamRecords read the schema from the stream instead of
taking it per call, matching Pipelines.Stream. The schema is a client-side
codec and is never persisted, so changing it never updates the stream.
…bles

Pipelines covers stream -> SQL -> sink (Iceberg and R2-file sinks, producers,
the HTTP endpoint default). Iceberg tables covers the catalog, Basin.Table,
maintenance, delete policies, and querying from other engines.
- R2.DataCatalog / Basin.Catalog: bucket: string | R2.Bucket (bucketName
  kept as a deprecated alias; switching forms is not a change)
- K2.Subscription: stream: string | K2.Stream (was streamId)
- Basin.Table: catalog: string | Basin.Catalog (a string is the bucket name)
- Pipelines.Sink: catalog: string | Basin.Catalog; config.bucket:
  string | R2.Bucket
@sam-goodwin
sam-goodwin merged commit 992d49d into main Oct 6, 2026
8 checks passed
@sam-goodwin
sam-goodwin deleted the feat/cloudflare-k2-basin branch October 6, 2026 23:18
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant