Documentation
¶
Index ¶
- Constants
- Variables
- func MarshallPayloadToJson() middleware.Middleware
- func MetadataAsAttributes(msg *message.Message) map[string]string
- func UnmarshallPayloadFromJson[T any](payloadType T) middleware.Middleware
- type AttributesProvider
- type MessageUnmarshaller
- type OrderingKeyProvider
- type Publisher
- type PublisherOption
- type PublisherOptions
- type Subscriber
- type SubscriberOption
- type SubscriberOptions
- type TopicPublisher
Constants ¶
const DefaultProcessingTimeout = 600 * time.Second
Variables ¶
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 UnmarshallPayloadFromJson ¶
func UnmarshallPayloadFromJson[T any](payloadType T) middleware.Middleware
Types ¶
type MessageUnmarshaller ¶
type OrderingKeyProvider ¶
type Publisher ¶
type Publisher struct {
// contains filtered or unexported fields
}
func NewGooglePublisher ¶
func NewGooglePublisher( c *pubsub.Client, routingFunc publisher.RoutingFunc, opts ...PublisherOption) (*Publisher, 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