Documentation
¶
Index ¶
- Variables
- func CreatePoller(ms MessageStore, worker SubscriptionWorker, config *SubscriberConfig) (*poller, error)
- func NewID() uuid.UUID
- func Pack(source interface{}) (map[string]interface{}, error)
- func Unpack(source map[string]interface{}, dest interface{}) error
- type Command
- type Event
- type GetOption
- func BatchSize(batchsize int) GetOption
- func Category(category string) GetOption
- func CommandStream(category string) GetOption
- func Converter(converter MessageConverter) GetOption
- func EventStream(category string, entityID uuid.UUID) GetOption
- func Last() GetOption
- func PositionStream(subscriberID string) GetOption
- func SincePosition(position int64) GetOption
- func SinceVersion(version int64) GetOption
- type Message
- type MessageConverter
- type MessageHandler
- type MessageReducer
- type MessageReducerConfig
- type MessageStore
- type Poller
- type Projector
- type ProjectorOption
- type Subscriber
- type SubscriberConfig
- type SubscriberOption
- func PollErrorDelay(pollErrorDelay time.Duration) SubscriberOption
- func PollTime(pollTime time.Duration) SubscriberOption
- func SubscribeBatchSize(batchSize int) SubscriberOption
- func SubscribeToCategory(category string) SubscriberOption
- func SubscribeToCommandStream(category string) SubscriberOption
- func SubscribeToEntityStream(category string, entityID uuid.UUID) SubscriberOption
- func UpdatePositionEvery(msgInterval int) SubscriberOption
- type SubscriptionWorker
- type WriteOption
Constants ¶
This section is empty.
Variables ¶
var ( ErrInvalidOptionCombination = errors.New("Cannot have the current combination of options for Get()") ErrSubscriberCannotUseBothStreamAndCategory = errors.New("Subscriber function cannot use both Stream and Category") ErrInvalidPollTime = errors.New("Invalid Subscriber poll time provided, can not be negative or zero") ErrInvalidPollErrorDelay = errors.New("Invalid Subscriber poll error delay provided, can not be negative or zero") ErrInvalidBatchSize = errors.New("Invalid Subscriber batch size provided, can not be negative or zero") ErrInvalidMsgInterval = errors.New("MsgInterval cannot be less than 2") ErrSubscriberNeedsCategoryOrStream = errors.New("Subscriber needs at least one of category or stream to be set upon creation") ErrSubscriberIDCannotBeEmpty = errors.New("Subscriber ID cannot be nil") ErrSubscriberIDCannotContainPlus = errors.New("Subscriber ID cannot contain a plus") ErrSubscriberIDCannotContainHyphen = errors.New("Subscriber ID cannot contain a hyphen") ErrPositionVersionMissing = errors.New("Subscriber cannot use non-versioned positionMsg object") ErrSubscriberNeedsAtLeastOneMessageHandler = errors.New("Subscriber needs at least one handler upon creation") ErrSubscriberCannotSubscribeToMultipleStreams = errors.New("Subscribers can only subscribe to one stream") ErrSubscriberCannotSubscribeToMultipleCategories = errors.New("Subscribers can only subscribe to one category") ErrProjectorNeedsAtLeastOneReducer = errors.New("Projector needs at least one reducer upon creation") ErrSubscriberMessageHandlerEqualToNil = errors.New("Subscriber Message Handler cannot be equal to nil") ErrSubscriberMessageHandlersEqualToNil = errors.New("Subscriber Message Handler array cannot be equal to nil") ErrSubscriberNilOption = errors.New("Options cannot include an option whose value is equal to nil") ErrDefaultStateNotSet = errors.New("Default state not set while trying to create a new projector") ErrDefaultStateCannotBePointer = errors.New("Default state cannot be a pointer when creating a projector") ErrGetMessagesCannotUseBothStreamAndCategory = errors.New("Get messages function cannot use both Stream and Category") ErrMessageNoID = errors.New("Message cannot be written without a new UUID") ErrGetMessagesRequiresEitherStreamOrCategory = errors.New("Get messages function must have either Stream or Category") ErrGetLastRequiresStream = errors.New("Get Last option requires a stream") ErrIncorrectNumberOfPositionsFound = errors.New("Exactly one position should be found per subscriber") ErrInvalidHandler = errors.New("Handler cannot be nil") ErrIncorrectMessageInPositionStream = errors.New("Position streams can only have position messages") ErrHandlerError = errors.New("Handler failed to handle message") ErrMissingMessageType = errors.New("All messages require a type") ErrMissingMessageCategory = errors.New("All messages require a category") ErrInvalidMessageCategory = errors.New("Hyphens are not allowed in category names") ErrInvalidCommandStream = errors.New("Hyphens are not allowed in command stream name") ErrInvalidEventStream = errors.New("Hyphens are not allowed in event stream name") ErrInvalidSubscriberID = errors.New("Hyphens and plusses are not allowed in subscriber ID") ErrInvalidPositionStream = errors.New("Position stream expects to have a single plus diving the subscriber ID and the word 'position'") ErrMissingMessageCategoryID = errors.New("All messages require a category ID") ErrMissingMessageData = errors.New("Messages payload must not be nil") ErrUnserializableData = errors.New("Message data could not be encoded as json") ErrDataIsNilPointer = errors.New("Message data is a nil pointer") ErrMissingGetOptions = errors.New("Options are required for the Get command") )
Errors
var NilUUID = uuid.Nil
So users don't have to import gomessagestore/uuid
Functions ¶
func CreatePoller ¶
func CreatePoller(ms MessageStore, worker SubscriptionWorker, config *SubscriberConfig) (*poller, error)
Types ¶
type Command ¶
type Command struct {
ID uuid.UUID
StreamCategory string
MessageType string
MessageVersion int64
GlobalPosition int64
Data map[string]interface{}
Metadata map[string]interface{}
Time time.Time
}
Command the model for writing a command to the Message Store
func (*Command) ToEnvelope ¶
func (cmd *Command) ToEnvelope() (*repository.MessageEnvelope, error)
ToEnvelope Allows for exporting to a MessageEnvelope type.
type Event ¶
type Event struct {
ID uuid.UUID
EntityID uuid.UUID
StreamCategory string
MessageType string
MessageVersion int64
GlobalPosition int64
Data map[string]interface{}
Metadata map[string]interface{}
Time time.Time
}
Event the model for writing an event to the Message Store
func (*Event) ToEnvelope ¶
func (event *Event) ToEnvelope() (*repository.MessageEnvelope, error)
ToEnvelope Allows for exporting to a MessageEnvelope type.
type GetOption ¶
type GetOption func(g *getOpts) error
GetOption provide optional arguments to the Get function
func CommandStream ¶
CommandStream allows for writing messages using an expected position
func Converter ¶
func Converter(converter MessageConverter) GetOption
Converter allows for automatic converting of non-Command/Event type messages
func EventStream ¶
EventStream allows for getting events in a specific stream
func Last ¶
func Last() GetOption
Last allows for getting only the most recent message (still returns an array)
func PositionStream ¶
Position allows for getting messages by position subscriber
func SincePosition ¶
SincePosition allows for getting only more recent messages
func SinceVersion ¶
SinceVersion allows for getting only more recent messages
type Message ¶
type Message interface {
ToEnvelope() (*repository.MessageEnvelope, error)
Type() string
Version() int64
Position() int64
}
Message Defines an interface that can define more specific class of messages such as Commands or Events.
func MsgEnvelopesToMessages ¶
func MsgEnvelopesToMessages(msgEnvelopes []*repository.MessageEnvelope, converters ...MessageConverter) []Message
MsgEnvelopesToMessages converts envelopes to any number of different structs that impliment the Message interface
type MessageConverter ¶
type MessageConverter func(*repository.MessageEnvelope) (Message, error)
MessageConverter allows the MsgEnvelopesToMessages to convert to structs that aren't defined in this library if the message isn't the correct type, returning nil or an error will cause the default converters to run
type MessageHandler ¶
type MessageReducer ¶
type MessageReducer interface {
Reduce(msg Message, previousState interface{}) interface{}
Type() string
}
MessageReducer Defines the expected behaviours of a reducer that ultimately is used by the projectors.
type MessageReducerConfig ¶
type MessageReducerConfig struct {
Reducer MessageReducer
Type string
}
MessageReducerConfig Contains all of the information needed to use a given reducer.
type MessageStore ¶
type MessageStore interface {
Write(ctx context.Context, message Message, opts ...WriteOption) error
Get(ctx context.Context, opts ...GetOption) ([]Message, error)
CreateProjector(opts ...ProjectorOption) (Projector, error)
CreateSubscriber(subscriberID string, handlers []MessageHandler, opts ...SubscriberOption) (Subscriber, error)
}
MessageStore Establishes the interface for Eventide.
func NewMessageStore ¶
func NewMessageStore(injectedDB *sql.DB) MessageStore
NewMessageStore Grabs a MessageStore instance.
func NewMessageStoreFromRepository ¶
func NewMessageStoreFromRepository(injectedRepo repository.Repository) MessageStore
NewMessageStoreFromRepository Grabs a MessageStore instance.
func NewMockMessageStoreWithMessages ¶ added in v0.2.2
func NewMockMessageStoreWithMessages(msgs []Message) MessageStore
NewMockMessageStoreWithMessages
type Projector ¶
type Projector interface {
Run(ctx context.Context, category string, entityID uuid.UUID) (interface{}, error)
}
Projector A base level interface that defines the projection functionality of gomessagestore.
type ProjectorOption ¶
type ProjectorOption func(proj *projector)
ReducerOption Variadic parameter support for reducers.
func DefaultState ¶
func DefaultState(defaultState interface{}) ProjectorOption
DefaultState registers a default state for use with a projector
func WithReducer ¶
func WithReducer(reducer MessageReducer) ProjectorOption
WithReducer registers a ruducer with the new projector
type Subscriber ¶
Subscriber allows for reaching out to the message service on a continual basis
func CreateSubscriberWithPoller ¶
func CreateSubscriberWithPoller(ms MessageStore, subscriberID string, handlers []MessageHandler, poller Poller, opts ...SubscriberOption) (Subscriber, error)
CreateSubscriberWithPoller is used for testing with dependency injection
type SubscriberConfig ¶
type SubscriberConfig struct {
// contains filtered or unexported fields
}
func GetSubscriberConfig ¶
func GetSubscriberConfig(opts ...SubscriberOption) (*SubscriberConfig, error)
GetSubscriberConfig changes SubscriberOptions into a valid SubscriberConfig object, or returns an error
type SubscriberOption ¶
type SubscriberOption func(config *SubscriberConfig) error
SubscriberOption allows for various options when creating a subscriber
func PollErrorDelay ¶
func PollErrorDelay(pollErrorDelay time.Duration) SubscriberOption
PollErrorDelay sets the interval between handling operations when Poll() errors
func PollTime ¶
func PollTime(pollTime time.Duration) SubscriberOption
PollTime sets the interval between handling operations
func SubscribeBatchSize ¶
func SubscribeBatchSize(batchSize int) SubscriberOption
SubscribeBatchSize sets the amount of messages to retrieve in a single handling operation
func SubscribeToCategory ¶
func SubscribeToCategory(category string) SubscriberOption
Subscribe to a category of streams
func SubscribeToCommandStream ¶
func SubscribeToCommandStream(category string) SubscriberOption
Subscribe to a specific command stream
func SubscribeToEntityStream ¶
func SubscribeToEntityStream(category string, entityID uuid.UUID) SubscriberOption
Subscribe to a specific entity stream
func UpdatePositionEvery ¶
func UpdatePositionEvery(msgInterval int) SubscriberOption
UpdatePostionEvery updates position of subscriber based on a msgInterval (cannot be < 2) An interval of 1 would create an event on every message, and possibly be picked up by itself, creating another event, and so on
type SubscriptionWorker ¶
type SubscriptionWorker interface {
GetMessages(ctx context.Context, position int64) ([]Message, error)
ProcessMessages(ctx context.Context, msgs []Message) (messagesHandled int, positionOfLastHandled int64, err error)
GetPosition(ctx context.Context) (int64, error)
SetPosition(ctx context.Context, position int64) error
}
func CreateWorker ¶
func CreateWorker(ms MessageStore, subscriberID string, handlers []MessageHandler, config *SubscriberConfig) (SubscriptionWorker, error)
type WriteOption ¶
type WriteOption func(w *writer)
WriteOption provide optional arguments to the Write function
func AtPosition ¶
func AtPosition(position int64) WriteOption
AtPosition allows for writing messages using an expected position