Skip to content

GCP.PubSub reference

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.

const acknowledge = yield* GCP.PubSub.Acknowledge(subscription);
yield* acknowledge({
body: { ackIds: [ackId] },
});

Source: src/GCP/PubSub/AcknowledgeHttp.ts Kind: Layer · Provides: GCP.PubSub.Acknowledge

HTTP implementation of Acknowledge.

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.

const getSchema = yield* GCP.PubSub.GetSchema(schema);
const live = yield* getSchema({ view: "FULL" });

Source: src/GCP/PubSub/GetSchemaHttp.ts Kind: Layer · Provides: GCP.PubSub.GetSchema

HTTP implementation of GetSchema.

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.

const publish = yield* GCP.PubSub.Publish(topic);
const { messageIds } = yield* publish({
body: { messages: [{ data: btoa("hello") }] },
});

Source: src/GCP/PubSub/PublishHttp.ts Kind: 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.

const pull = yield* GCP.PubSub.Pull(subscription);
const { receivedMessages } = yield* pull({
body: { maxMessages: 10 },
});

Source: src/GCP/PubSub/PullHttp.ts Kind: Layer · Provides: GCP.PubSub.Pull

HTTP implementation of Pull.

Source: src/GCP/PubSub/ReadSubscription.ts

Consume access to a Pub/Sub Subscription: pull, acknowledge, modifyAckDeadline. Grants roles/pubsub.subscriber on the subscription only.

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

Source: src/GCP/PubSub/ReadSubscriptionHttp.ts Kind: Layer · Provides: GCP.PubSub.ReadSubscription

HTTP implementation of ReadSubscription over the Pub/Sub REST API.

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.

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

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

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

Source: src/GCP/PubSub/Subscription.ts

A Google Cloud Pub/Sub 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" },
});
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.

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" },
});
const topic = yield* GCP.PubSub.Topic("events", {
labels: { env: "prod" },
});
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])),
) {}

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 account roles/run.invoker on 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, grants roles/pubsub.subscriber on it, and acks each batch after the handler succeeds.

Consume it through consumeTopicMessages.

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)

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.

const validate = yield* GCP.PubSub.ValidateMessage(schema);
yield* validate({
encoding: "JSON",
message: btoa(JSON.stringify({ id: "abc" })),
});

Source: src/GCP/PubSub/ValidateMessageHttp.ts Kind: Layer · Provides: GCP.PubSub.ValidateMessage

HTTP implementation of ValidateMessage.

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.

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

Source: src/GCP/PubSub/WriteTopicHttp.ts Kind: Layer · Provides: GCP.PubSub.WriteTopic

HTTP implementation of WriteTopic over the Pub/Sub REST API.