Repository navigation
feat(cloudflare)!: K2 streams and Basin (Pipelines, Catalog, Tables) - #2010
Merged
Merged
Conversation
Contributor
|
Install the packages built from this commit: alchemypnpm install https://pkg.alchemy.run/alchemy/pr:2010:f5e2556@alchemy.run (6)pnpm install https://pkg.alchemy.run/@alchemy.run/better-auth/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@alchemy.run/cloudflare-runtime/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@alchemy.run/frontend-frameworks/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@alchemy.run/node-utils/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@alchemy.run/floci/pr:2010:f5e2556pnpm 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:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/acme/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/aws/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/axiom/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/cloudflare/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/doppler/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/fly-io/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/github/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/hetzner/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/infisical/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/neon/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/prisma/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/planetscale/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/railway/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/stripe/pr:2010:f5e2556pnpm install https://pkg.alchemy.run/@distilled.cloud/zerossl/pr:2010:f5e2556Published |
sam-goodwin
force-pushed
the
feat/cloudflare-k2-basin
branch
3 times, most recently
from
October 6, 2026 04:37
b3bae4b to
e6ebcce
Compare
sam-goodwin
marked this pull request as ready for review
October 6, 2026 08:33
sam-goodwin
force-pushed
the
feat/cloudflare-k2-basin
branch
2 times, most recently
from
October 6, 2026 11:04
250d4f6 to
b0d2fc6
Compare
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
force-pushed
the
feat/cloudflare-k2-basin
branch
from
October 6, 2026 16:24
b0d2fc6 to
44a6d9b
Compare
…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.
This was referenced Oct 6, 2026
…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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Cloudflare K2 streams and a
Cloudflare.Basinnamespace (Pipelines + Catalog + Iceberg tables), with hand-written guides for both. Builds on alchemy-run/distilled#730 (merged), whichmain's distilled pin already includes.K2: a durable log you consume yourself
The schema lives on the stream, so every binding on it is typed the same way:
Produce from an Effect Worker:
Swap the layer to produce from somewhere else:
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:
Each consuming host gets its own subscription (fan-out). To split records between hosts, pass a shared
Subscription:Or drive the leases yourself:
Async Workers get a typed binding:
Basin: Pipelines + Catalog + Iceberg tables
Cloudflare.Basin.*holds the same values asCloudflare.Pipelines.*andCloudflare.R2.DataCatalog. The old paths and persisted resource types are unchanged.Producers work like K2's and are typed by the stream's schema:
Tables you own directly go over Iceberg REST, using
@distilled.cloud/iceberg: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'sbucketNamestill works as a deprecated alias.Docs
/cloudflare/messaging/k2)./cloudflare/data/pipelines): stream → SQL → Iceberg or R2-file sink./cloudflare/data/iceberg-tables): the catalog,Basin.Table, and querying from other engines.Dependencies
@distilled.cloud/iceberg.typescript, so the lockfile keeps resolving against the workspace TypeScript.