google

package module
v0.0.0-...-7593c98 Latest Latest
Warning

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

Go to latest
Published: Aug 27, 2025 License: MIT Imports: 11 Imported by: 0

Documentation

Index

Constants

View Source
const DefaultProcessingTimeout = 600 * time.Second

Variables

This section is empty.

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 PubsubTopic

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

func (PubsubTopic) Close

func (p PubsubTopic) Close() error

func (PubsubTopic) GetMessageID

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

func (PubsubTopic) Publish

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

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

func (s *Subscriber) Close() error

func (*Subscriber) Subscribe

func (s *Subscriber) Subscribe(ctx context.Context) (<-chan *message.Message, error)

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 will dictate how long a message will be processed before it is nacked. 0 means no timeout, wait forever. Keep in mind that GCP uses the "Acknowledgement deadline" to determine if a message needs to be redelivered. ProcessingTimeout has no impact on the "Acknowledgement deadline". Default value is 600 seconds, which is the max value of the GCP "Acknowledgement deadline".

func WithProcessingTimeoutHandler

func WithProcessingTimeoutHandler(handler func(ctx context.Context, msg *pubsub.Message)) SubscriberOption

WithProcessingTimeoutHandler is a function that is called when a message processing times out.

func WithReceiveSettings

func WithReceiveSettings(settings pubsub.ReceiveSettings) SubscriberOption

WithReceiveSettings is a set of options to pass the underlying gcp pubsub.Subscription

type SubscriberOptions

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

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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