Documentation ¶
Index ¶
- Constants
- func ApplyMetadata(cloudEvent map[string]interface{}, componentFeatures []Feature, ...)
- func ConvertTLSPropertiesToTLSConfig(properties TLSProperties) (*tls.Config, error)
- func FromCloudEvent(cloudEvent []byte, topic, pubsub, traceParent string, traceState string) (map[string]interface{}, error)
- func FromRawPayload(data []byte, topic, pubsub string) map[string]interface{}
- func HasExpired(cloudEvent map[string]interface{}) bool
- func NewCloudEventsEnvelope(id, source, eventType, subject string, topic string, pubsubName string, ...) map[string]interface{}
- func Ping(ctx context.Context, pubsub PubSub) error
- type AppBulkResponse
- type AppBulkResponseEntry
- type AppResponse
- type AppResponseStatus
- type BulkHandler
- type BulkMessage
- type BulkMessageEntry
- type BulkPublishRequest
- type BulkPublishResponse
- type BulkPublishResponseFailedEntry
- type BulkPublisher
- type BulkSubscribeConfig
- type BulkSubscribeResponse
- type BulkSubscribeResponseEntry
- type BulkSubscriber
- type ConcurrencyMode
- type Feature
- type Handler
- type Metadata
- type NewMessage
- type PubSub
- type PublishRequest
- type SubscribeRequest
- type TLSProperties
Constants ¶
const ( // DefaultCloudEventType is the default event type for an Dapr published event. DefaultCloudEventType = "com.dapr.event.sent" // DefaultBulkEventType is the default bulk event type for a Dapr published event. DefaultBulkEventType = "com.dapr.event.sent.bulk" // CloudEventsSpecVersion is the specversion used by Dapr for the cloud events implementation. CloudEventsSpecVersion = "1.0" // DefaultCloudEventSource is the default event source. DefaultCloudEventSource = "Dapr" // DefaultCloudEventDataContentType is the default content-type for the data attribute. DefaultCloudEventDataContentType = "text/plain" // traceid, backwards compatibles. // ::TODO delete traceid, and keep traceparent. TraceIDField = "traceid" TraceParentField = "traceparent" TraceStateField = "tracestate" TopicField = "topic" PubsubField = "pubsubname" ExpirationField = "expiration" DataContentTypeField = "datacontenttype" DataField = "data" DataBase64Field = "data_base64" SpecVersionField = "specversion" TypeField = "type" SourceField = "source" IDField = "id" SubjectField = "subject" TimeField = "time" )
const ( // CACert is the metadata key name for the CA certificate. CACert = "caCert" // ClientCert is the metadata key name for the client certificate. ClientCert = "clientCert" // ClientKey is the metadata key name for the client key. ClientKey = "clientKey" )
const RuntimeConsumerIDKey = "consumerID"
When the Dapr component does not explicitly specify a consumer group, this value provided by the runtime must be used. This value is specific to each Dapr App. As a result, by default, each Dapr App will receive all messages published to the topic at least once. See https://github.com/dapr/dapr/blob/21566de8d7fdc7d43ae627ffc0698cc073fa71b0/pkg/runtime/runtime.go#L1735-L1739
Variables ¶
This section is empty.
Functions ¶
func ApplyMetadata ¶
func ApplyMetadata(cloudEvent map[string]interface{}, componentFeatures []Feature, metadata map[string]string)
ApplyMetadata will process metadata to modify the cloud event based on the component's feature set.
func ConvertTLSPropertiesToTLSConfig ¶
func ConvertTLSPropertiesToTLSConfig(properties TLSProperties) (*tls.Config, error)
ConvertTLSPropertiesToTLSConfig converts the TLSProperties to a tls.Config.
func FromCloudEvent ¶
func FromCloudEvent(cloudEvent []byte, topic, pubsub, traceParent string, traceState string) (map[string]interface{}, error)
FromCloudEvent returns a map representation of an existing cloudevents JSON.
func FromRawPayload ¶
FromRawPayload returns a CloudEvent for a raw payload on subscriber's end.
func HasExpired ¶
HasExpired determines if the current cloud event has expired.
func NewCloudEventsEnvelope ¶
func NewCloudEventsEnvelope(id, source, eventType, subject string, topic string, pubsubName string, dataContentType string, data []byte, traceParent string, traceState string, ) map[string]interface{}
NewCloudEventsEnvelope returns a map representation of a cloudevents JSON.
Types ¶
type AppBulkResponse ¶
type AppBulkResponse struct {
AppResponses []AppBulkResponseEntry `json:"statuses"`
}
AppBulkResponse is the whole bulk subscribe response sent by App
type AppBulkResponseEntry ¶
type AppBulkResponseEntry struct { EntryId string `json:"entryId"` //nolint:stylecheck Status AppResponseStatus `json:"status"` }
AppBulkResponseEntry Represents single response, as part of AppBulkResponse, to be sent by subscibed App for the corresponding single message during bulk subscribe
type AppResponse ¶
type AppResponse struct {
Status AppResponseStatus `json:"status"`
}
AppResponse is the object describing the response from user code after a pubsub event.
type AppResponseStatus ¶
type AppResponseStatus string
AppResponseStatus represents a status of a PubSub response.
const ( // Success means the message is received and processed correctly. Success AppResponseStatus = "SUCCESS" // Retry means the message is received but could not be processed and must be retried. Retry AppResponseStatus = "RETRY" // Drop means the message is received but should not be processed. Drop AppResponseStatus = "DROP" )
type BulkHandler ¶
type BulkHandler func(ctx context.Context, msg *BulkMessage) ([]BulkSubscribeResponseEntry, error)
[]BulkSubscribeResponseEntry represents individual statuses for each message in an orderly fashion.
type BulkMessage ¶
type BulkMessage struct { Entries []BulkMessageEntry `json:"entries"` Topic string `json:"topic"` Metadata map[string]string `json:"metadata"` }
BulkMessage represents bulk message arriving from a message bus instance.
func (BulkMessage) String ¶
func (m BulkMessage) String() string
String implements fmt.Stringer and it's useful for debugging.
type BulkMessageEntry ¶
type BulkMessageEntry struct { EntryId string `json:"entryId"` //nolint:stylecheck Event []byte `json:"event"` ContentType string `json:"contentType,omitempty"` Metadata map[string]string `json:"metadata"` }
BulkMessageEntry represents a single message inside a bulk request.
func (BulkMessageEntry) String ¶
func (m BulkMessageEntry) String() string
String implements fmt.Stringer and it's useful for debugging.
type BulkPublishRequest ¶
type BulkPublishRequest struct { Entries []BulkMessageEntry `json:"entries"` PubsubName string `json:"pubsubname"` Topic string `json:"topic"` Metadata map[string]string `json:"metadata"` }
BulkPublishRequest is the request to publish mutilple messages.
type BulkPublishResponse ¶
type BulkPublishResponse struct {
FailedEntries []BulkPublishResponseFailedEntry `json:"failedEntries"`
}
BulkPublishResponse contains the list of failed entries in a bulk publish request.
func NewBulkPublishResponse ¶
func NewBulkPublishResponse(messages []BulkMessageEntry, err error) BulkPublishResponse
NewBulkPublishResponse returns a BulkPublishResponse with each entry having same error. This method is a helper method to map a single error response on BulkPublish to multiple events.
type BulkPublishResponseFailedEntry ¶
type BulkPublishResponseFailedEntry struct { EntryId string `json:"entryId"` //nolint:stylecheck Error error `json:"error"` }
BulkPublishResponseFailedEntry Represents single publish response, as part of BulkPublishResponse to be sent to publishing App for the corresponding single message during bulk publish
type BulkPublisher ¶
type BulkPublisher interface {
BulkPublish(ctx context.Context, req *BulkPublishRequest) (BulkPublishResponse, error)
}
BulkPublish publishes a collection of entries/messages in a BulkPublishRequest to a message bus topic and returns a BulkPublishResponse with failed entries for any failed messages. Error is returned on partial or complete failures. If there are no failures, error is nil.
type BulkSubscribeConfig ¶
type BulkSubscribeConfig struct { MaxMessagesCount int `json:"maxMessagesCount,omitempty"` MaxAwaitDurationMs int `json:"maxAwaitDurationMs,omitempty"` }
BulkSubscribeConfig represents the configuration for bulk subscribe. It depends on specific componets to support these.
type BulkSubscribeResponse ¶
type BulkSubscribeResponse struct { Error error `json:"error"` Statuses []BulkSubscribeResponseEntry `json:"statuses"` }
BulkSubscribeResponse is the whole bulk subscribe response sent to building block
type BulkSubscribeResponseEntry ¶
type BulkSubscribeResponseEntry struct { EntryId string `json:"entryId"` //nolint:stylecheck Error error `json:"error"` }
BulkSubscribeResponseEntry Represents single subscribe response item, as part of BulkSubscribeResponse to be sent to building block for the corresponding single message during bulk subscribe
type BulkSubscriber ¶
type BulkSubscriber interface { // BulkSubscribe is used to subscribe to a topic and receive collection of entries/ messages // from a message bus topic. // The bulkHandler will be called with a list of messages. BulkSubscribe(ctx context.Context, req SubscribeRequest, bulkHandler BulkHandler) error }
BulkSubscriber is the interface defining BulkSubscribe definition for message buses
type ConcurrencyMode ¶
type ConcurrencyMode string
ConcurrencyMode is a pub/sub metadata setting that allows to specify whether messages are delivered in a serial or parallel execution.
const ( // ConcurrencyKey is the metadata key name for ConcurrencyMode. ConcurrencyKey = "concurrencyMode" Single ConcurrencyMode = "single" Parallel ConcurrencyMode = "parallel" )
func Concurrency ¶
func Concurrency(metadata map[string]string) (ConcurrencyMode, error)
Concurrency takes a metadata object and returns the ConcurrencyMode configured. Default is Parallel.
type Feature ¶
type Feature string
Feature names a feature that can be implemented by PubSub components.
type Handler ¶
type Handler func(ctx context.Context, msg *NewMessage) error
Handler is the handler used to invoke the app handler.
type NewMessage ¶
type NewMessage struct { Data []byte `json:"data"` Topic string `json:"topic"` Metadata map[string]string `json:"metadata"` ContentType *string `json:"contentType,omitempty"` }
NewMessage is an event arriving from a message bus instance.
func (NewMessage) String ¶
func (m NewMessage) String() string
String implements fmt.Stringer and it's useful for debugging.
type PubSub ¶
type PubSub interface { Init(ctx context.Context, metadata Metadata) error Features() []Feature Publish(ctx context.Context, req *PublishRequest) error Subscribe(ctx context.Context, req SubscribeRequest, handler Handler) error Close() error GetComponentMetadata() map[string]string }
PubSub is the interface for message buses.
type PublishRequest ¶
type PublishRequest struct { Data []byte `json:"data"` PubsubName string `json:"pubsubname"` Topic string `json:"topic"` Metadata map[string]string `json:"metadata"` ContentType *string `json:"contentType,omitempty"` }
PublishRequest is the request to publish a message.
type SubscribeRequest ¶
type SubscribeRequest struct { Topic string `json:"topic"` Metadata map[string]string `json:"metadata"` BulkSubscribeConfig BulkSubscribeConfig `json:"bulkSubscribe,omitempty"` }
SubscribeRequest is the request to subscribe to a topic.
type TLSProperties ¶
TLSProperties is a struct that contains the TLS properties.
Source Files ¶
Directories ¶
Path | Synopsis |
---|---|
aws
|
|
azure
|
|
gcp
|
|
Package natsstreaming implements NATS Streaming pubsub component
|
Package natsstreaming implements NATS Streaming pubsub component |
solace
|
|