google

package module
v0.0.0-...-f414a8c Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 1, 2026 License: MIT Imports: 13 Imported by: 0

Documentation

Index

Constants

View Source
const DefaultProcessingTimeout = 600 * time.Second

Variables

View Source
var ErrSubscriptionNotFound = errors.New("subscription does not exist")

ErrSubscriptionNotFound is returned by Receive when the subscription does not exist.

Functions

func MarshallPayloadToJson

func MarshallPayloadToJson() middleware.Middleware

func MetadataAsAttributes

func MetadataAsAttributes(msg *message.Message) map[string]string

func UnmarshallPayloadFromJson

func UnmarshallPayloadFromJson[T any](payloadType T) middleware.Middleware

Types

type AttributesProvider

type AttributesProvider func(message *message.Message) map[string]string

type MessageUnmarshaller

type MessageUnmarshaller func(ctx context.Context, msg *pubsub.Message) (*message.Message, error)

type OrderingKeyProvider

type OrderingKeyProvider func(message *message.Message) string

type Publisher

type Publisher struct {
	// contains filtered or unexported fields
}

func NewGooglePublisher

func NewGooglePublisher(
	c *pubsub.Client,
	routingFunc publisher.RoutingFunc,
	opts ...PublisherOption) (*Publisher, error)

func (*Publisher) Close

func (p *Publisher) Close() error

func (*Publisher) Publish

func (p *Publisher) Publish(message *message.Message) error

type PublisherOption

type PublisherOption func(*PublisherOptions)

func WithAttributesProvider

func WithAttributesProvider(provider AttributesProvider) PublisherOption

WithAttributesProvider is a function that returns attributes for a given message. If not provided, no attributes are used. A provider to set the attribute on the pubsub message. By default, it's using MetadataAsAttributes which converts all metadata entries as attributes.

func WithOrderingKeyProvider

func WithOrderingKeyProvider(provider OrderingKeyProvider) PublisherOption

WithOrderingKeyProvider is a function that returns an ordering key for a given message. If not provided, no ordering key is used.

type PublisherOptions

type PublisherOptions struct {
	// contains filtered or unexported fields
}

type Subscriber

type Subscriber struct {
	// contains filtered or unexported fields
}

func NewGoogleSubscriber

func NewGoogleSubscriber(
	c *pubsub.Client,
	subscription string,
	opts ...SubscriberOption) (*Subscriber, error)

func (*Subscriber) Receive

func (s *Subscriber) Receive(ctx context.Context, handler message.HandlerFunc) error

Receive implements subscriber.Subscriber. The Pub/Sub client limits the concurrency, see WithReceiveSettings.

type SubscriberOption

type SubscriberOption func(*SubscriberOptions)

func WithParseAttributes

func WithParseAttributes(parseAttributes bool) SubscriberOption

WithParseAttributes is a flag to indicate if the attributes should be parsed or not, meaning that boolean true/false, integers and floats are going to be their respective types. The default is to just keep everything as strings.

func WithProcessingTimeout

func WithProcessingTimeout(timeout time.Duration) SubscriberOption

WithProcessingTimeout is the deadline of the context of each message, see subscriber.DispatchOptions. 0 means no deadline. The Pub/Sub client extends the "Acknowledgement deadline" while the handler runs, up to ReceiveSettings.MaxExtension. Default value is 600 seconds, which is the max value of the GCP "Acknowledgement deadline".

func WithReceiveSettings

func WithReceiveSettings(settings pubsub.ReceiveSettings) SubscriberOption

WithReceiveSettings is a set of options to pass the underlying gcp pubsub.Subscriber. Its MaxOutstandingMessages and MaxOutstandingBytes limit how many messages are handled at the same time. Leave ShutdownOptions nil to keep the shutdown behavior of Receive (wait for the in-flight handlers): with ShutdownOptions set, the Pub/Sub client may stop waiting for them, and nack them.

type SubscriberOptions

type SubscriberOptions struct {
	// contains filtered or unexported fields
}

type TopicPublisher

type TopicPublisher struct {
	// contains filtered or unexported fields
}

TopicPublisher publishes the messages of a single topic.

func (TopicPublisher) Close

func (p TopicPublisher) Close() error

func (TopicPublisher) GetMessageID

func (p TopicPublisher) GetMessageID(message *pubsub.Message) string

func (TopicPublisher) Publish

func (p TopicPublisher) Publish(ctx context.Context, message *pubsub.Message) error

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL