GCP.PubSub reference
Acknowledge
Section titled “Acknowledge”Source:
src/GCP/PubSub/Acknowledge.ts
Runtime binding for Pub/Sub subscriptions.acknowledge.
Bind this operation to a Subscription in a Function/Action init
phase. Provide AcknowledgeHttp.
Acknowledge: Acknowledging Messages
Section titled “Acknowledge: Acknowledging Messages”const acknowledge = yield* GCP.PubSub.Acknowledge(subscription);yield* acknowledge({ body: { ackIds: [ackId] },});AcknowledgeHttp
Section titled “AcknowledgeHttp”Source:
src/GCP/PubSub/AcknowledgeHttp.tsKind: Layer · Provides:GCP.PubSub.Acknowledge
HTTP implementation of Acknowledge.
GetSchema
Section titled “GetSchema”Source:
src/GCP/PubSub/GetSchema.ts
Runtime binding for Pub/Sub schemas.get.
Bind this operation to a Schema in a Function/Action init phase.
Provide GetSchemaHttp.
GetSchema: Reading a Schema
Section titled “GetSchema: Reading a Schema”const getSchema = yield* GCP.PubSub.GetSchema(schema);const live = yield* getSchema({ view: "FULL" });GetSchemaHttp
Section titled “GetSchemaHttp”Source:
src/GCP/PubSub/GetSchemaHttp.tsKind: Layer · Provides:GCP.PubSub.GetSchema
HTTP implementation of GetSchema.
Publish
Section titled “Publish”Source:
src/GCP/PubSub/Publish.ts
Runtime binding for Pub/Sub topics.publish.
Bind this operation to a Topic in a Function/Action init phase.
Provide PublishHttp.
Publish: Publishing Messages
Section titled “Publish: Publishing Messages”const publish = yield* GCP.PubSub.Publish(topic);const { messageIds } = yield* publish({ body: { messages: [{ data: btoa("hello") }] },});PublishHttp
Section titled “PublishHttp”Source:
src/GCP/PubSub/PublishHttp.tsKind: Layer · Provides:GCP.PubSub.Publish
HTTP implementation of Publish.
Source:
src/GCP/PubSub/Pull.ts
Runtime binding for Pub/Sub subscriptions.pull.
Bind this operation to a Subscription in a Function/Action init
phase. Provide PullHttp.
Pull: Pulling Messages
Section titled “Pull: Pulling Messages”const pull = yield* GCP.PubSub.Pull(subscription);const { receivedMessages } = yield* pull({ body: { maxMessages: 10 },});PullHttp
Section titled “PullHttp”Source:
src/GCP/PubSub/PullHttp.tsKind: Layer · Provides:GCP.PubSub.Pull
HTTP implementation of Pull.
ReadSubscription
Section titled “ReadSubscription”Source:
src/GCP/PubSub/ReadSubscription.ts
Consume access to a Pub/Sub Subscription: pull, acknowledge,
modifyAckDeadline. Grants roles/pubsub.subscriber on the subscription
only.
ReadSubscription: Consuming messages
Section titled “ReadSubscription: Consuming messages”Pull and acknowledge
const inbox = yield* GCP.PubSub.ReadSubscription(subscription);const messages = yield* inbox.pull({ maxMessages: 10 });for (const message of messages) { yield* Effect.log(message.text, message.attributes);}yield* inbox.acknowledge(messages.map((m) => m.ackId));// …provided with Effect.provide(GCP.PubSub.ReadSubscriptionHttp)Release a message for redelivery
yield* inbox.modifyAckDeadline([message.ackId], 0);ReadSubscriptionHttp
Section titled “ReadSubscriptionHttp”Source:
src/GCP/PubSub/ReadSubscriptionHttp.tsKind: Layer · Provides:GCP.PubSub.ReadSubscription
HTTP implementation of ReadSubscription over the Pub/Sub REST API.
Schema
Section titled “Schema”Source:
src/GCP/PubSub/Schema.ts
A Google Cloud Pub/Sub schema used to validate published messages.
Pub/Sub schemas have no labels field. Ownership is by resource name;
list returns every schema in the project so pnpm nuke:gcp can
clean leaked test schemas.
Schema: Creating a Schema
Section titled “Schema: Creating a Schema”Generated name (Avro)
const schema = yield* GCP.PubSub.Schema("Events", { type: "AVRO", definition: JSON.stringify({ type: "record", name: "Event", fields: [{ name: "id", type: "string" }], }),});Explicit id and Protocol Buffer definition
const schema = yield* GCP.PubSub.Schema("Events", { schemaId: "order-events", type: "PROTOCOL_BUFFER", definition: 'syntax = "proto3";\nmessage Event { string id = 1; }',});Schema: Updating a Schema
Section titled “Schema: Updating a Schema”Change definition on the same logical id; the schema keeps its name and
a new revision is committed.
const schema = yield* GCP.PubSub.Schema("Events", { type: "AVRO", definition: JSON.stringify({ type: "record", name: "Event", fields: [ { name: "id", type: "string" }, { name: "count", type: "int", default: 0 }, ], }),});Snapshot
Section titled “Snapshot”Source:
src/GCP/PubSub/Snapshot.ts
A Google Cloud Pub/Sub snapshot of a subscription’s message backlog.
Snapshots capture unacked messages at creation time (and subsequent
publishes to the topic) so a subscription can later Seek back to that
point. The source subscription is immutable — changing it replaces the
snapshot. Labels can be updated in place.
Snapshot: Creating a Snapshot
Section titled “Snapshot: Creating a Snapshot”Snapshot of a pull subscription
const topic = yield* GCP.PubSub.Topic("events", {});const subscription = yield* GCP.PubSub.Subscription("orders", { topic: topic.name,});const snapshot = yield* GCP.PubSub.Snapshot("checkpoint", { subscription: subscription.name,});Explicit id and labels
const topic = yield* GCP.PubSub.Topic("events", {});const subscription = yield* GCP.PubSub.Subscription("orders", { topic: topic.name,});const snapshot = yield* GCP.PubSub.Snapshot("checkpoint", { snapshotId: "order-checkpoint", subscription: subscription.name, labels: { env: "prod" },});Subscription
Section titled “Subscription”Source:
src/GCP/PubSub/Subscription.ts
A Google Cloud Pub/Sub subscription.
Subscription: Creating a Subscription
Section titled “Subscription: Creating a Subscription”Pull subscription on a topic
const topic = yield* GCP.PubSub.Topic("events", {});const subscription = yield* GCP.PubSub.Subscription("orders", { topic: topic.name,});Explicit id, labels, and ack deadline
const topic = yield* GCP.PubSub.Topic("events", {});const subscription = yield* GCP.PubSub.Subscription("orders", { subscriptionId: "order-events", topic: topic.name, ackDeadlineSeconds: 30, labels: { env: "prod" },});Subscription: Pulling Messages
Section titled “Subscription: Pulling Messages”const pull = yield* GCP.PubSub.Pull(subscription);const { receivedMessages } = yield* pull({ body: { maxMessages: 1 },});const ackIds = (receivedMessages ?? []) .map((message) => message.ackId) .filter((ackId): ackId is string => ackId !== undefined);if (ackIds.length > 0) { const acknowledge = yield* GCP.PubSub.Acknowledge(subscription); yield* acknowledge({ body: { ackIds } });}Source:
src/GCP/PubSub/Topic.ts
A Google Cloud Pub/Sub topic.
Topic: Creating a Topic
Section titled “Topic: Creating a Topic”Generated name
const topic = yield* GCP.PubSub.Topic("events", {});Explicit id and labels
const topic = yield* GCP.PubSub.Topic("events", { topicId: "order-events", labels: { env: "prod" },});Topic: Updating a Topic
Section titled “Topic: Updating a Topic”const topic = yield* GCP.PubSub.Topic("events", { labels: { env: "prod" },});Topic: Binding from a Function
Section titled “Topic: Binding from a Function”export class Api extends GCP.Function<Api>()( "Api", { main: import.meta.url }, Effect.gen(function* () { const topic = yield* GCP.PubSub.Topic("events", {}); const publish = yield* GCP.PubSub.Publish(topic); return { fetch: Effect.gen(function* () { yield* publish({ body: { messages: [{ data: btoa("hello") }] }, }).pipe(Effect.orDie); return HttpServerResponse.text("ok"); }), }; }).pipe(Effect.provide([GCP.PubSub.PublishHttp])),) {}Topic: Destroying a Topic
Section titled “Topic: Destroying a Topic”TopicEventSource
Section titled “TopicEventSource”Source:
src/GCP/PubSub/TopicEventSource.ts
Event source connecting a Pub/Sub Topic to the hosting compute.
The host-specific implementation layers are:
GCP.Run.TopicEventSource— push delivery to an HTTP host (GCP.Run.Service/GCP.Function,GCP.CloudFunctions.Function). Provisions a push subscription at the host’s URL, authenticated with an OIDC token for the host’s runtime service account, grants that accountroles/run.invokeron the host, and verifies the token on every delivery. A 2xx response acks; a failed handler returns 500 and Pub/Sub redelivers.GCP.Run.TopicPullEventSource— a pull loop for hosts without an inbound URL (GCP.Run.Job,GCP.Run.WorkerPool). Provisions a pull subscription, grantsroles/pubsub.subscriberon it, and acks each batch after the handler succeeds.
Consume it through consumeTopicMessages.
TopicEventSource: Consuming a Topic
Section titled “TopicEventSource: Consuming a Topic”Push to a Cloud Run service
export class Worker extends GCP.Function<Worker>()( "Worker", { main: import.meta.url }, Effect.gen(function* () { const orders = yield* GCP.PubSub.Topic("Orders", {}); yield* GCP.PubSub.consumeTopicMessages(orders, (messages) => messages.pipe( Stream.runForEach(({ message }) => Effect.log(atob(message.data ?? "")), ), ), ); }).pipe(Effect.provide(GCP.Run.TopicEventSource)),) {}Dead-letter messages the handler keeps failing on
const failed = yield* GCP.PubSub.Topic("FailedOrders", {});yield* GCP.PubSub.consumeTopicMessages( orders, { deadLetter: { topic: failed, maxDeliveryAttempts: 5 } }, (messages) => messages.pipe(Stream.runForEach(handle)),);Pull from a worker pool
yield* GCP.PubSub.consumeTopicMessages( orders, { maxMessages: 50 }, (messages) => messages.pipe(Stream.runForEach(handle)),);// …provided with Effect.provide(GCP.Run.TopicPullEventSource)ValidateMessage
Section titled “ValidateMessage”Source:
src/GCP/PubSub/ValidateMessage.ts
Runtime binding for Pub/Sub schemas.validateMessage.
Bind this operation to a Schema in a Function/Action init phase.
Provide ValidateMessageHttp.
Grants roles/pubsub.viewer on the project: schemas.validateMessage is
authorized on the parent project, and a grant on the schema’s own IAM
policy is not enough.
ValidateMessage: Validating Messages
Section titled “ValidateMessage: Validating Messages”const validate = yield* GCP.PubSub.ValidateMessage(schema);yield* validate({ encoding: "JSON", message: btoa(JSON.stringify({ id: "abc" })),});ValidateMessageHttp
Section titled “ValidateMessageHttp”Source:
src/GCP/PubSub/ValidateMessageHttp.tsKind: Layer · Provides:GCP.PubSub.ValidateMessage
HTTP implementation of ValidateMessage.
WriteTopic
Section titled “WriteTopic”Source:
src/GCP/PubSub/WriteTopic.ts
Publish access to a Pub/Sub Topic: publish, publishBatch.
Grants roles/pubsub.publisher on the topic only. Topics have no runtime
read — consume through a subscription with ReadSubscription.
WriteTopic: Publishing messages
Section titled “WriteTopic: Publishing messages”Publish one message
const events = yield* GCP.PubSub.WriteTopic(topic);const messageId = yield* events.publish({ data: JSON.stringify({ type: "signup" }), attributes: { source: "api" },});// …provided with Effect.provide(GCP.PubSub.WriteTopicHttp)Publish a batch
const events = yield* GCP.PubSub.WriteTopic(topic);const ids = yield* events.publishBatch([{ data: "a" }, { data: "b" }]);WriteTopicHttp
Section titled “WriteTopicHttp”Source:
src/GCP/PubSub/WriteTopicHttp.tsKind: Layer · Provides:GCP.PubSub.WriteTopic
HTTP implementation of WriteTopic over the Pub/Sub REST API.