gomessagestore

package module
v0.3.0 Latest Latest
Warning

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

Go to latest
Published: Jul 18, 2019 License: MIT Imports: 12 Imported by: 0

README

gomessagestore

basic eventide interface for go

Documentation

Index

Constants

This section is empty.

Variables

View Source
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

View Source
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)

func NewID

func NewID() uuid.UUID

NewID creates a new UUID

func Pack

func Pack(source interface{}) (map[string]interface{}, error)

Pack packs a GO object into JSON-esque objects used in the Command and Event objects

func Unpack

func Unpack(source map[string]interface{}, dest interface{}) error

Unpack unpacks JSON-esque objects used in the Command and Event objects into GO objects

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

func (cmd *Command) Position() int64

Position gets the command's Position field

func (*Command) ToEnvelope

func (cmd *Command) ToEnvelope() (*repository.MessageEnvelope, error)

ToEnvelope Allows for exporting to a MessageEnvelope type.

func (*Command) Type

func (cmd *Command) Type() string

Type returns the type of this message (business action taking place)

func (*Command) Version

func (cmd *Command) Version() int64

Version gets the command's Version field

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

func (event *Event) Position() int64

Position gets the events's Position field

func (*Event) ToEnvelope

func (event *Event) ToEnvelope() (*repository.MessageEnvelope, error)

ToEnvelope Allows for exporting to a MessageEnvelope type.

func (*Event) Type

func (event *Event) Type() string

Type returns the type of this message (business action taking place)

func (*Event) Version

func (event *Event) Version() int64

Version gets the event's Version field

type GetOption

type GetOption func(g *getOpts) error

GetOption provide optional arguments to the Get function

func BatchSize

func BatchSize(batchsize int) GetOption

BatchSize changes how many messages are returned (default 1000)

func Category

func Category(category string) GetOption

Category allows for getting messages by category

func CommandStream

func CommandStream(category string) GetOption

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

func EventStream(category string, entityID uuid.UUID) GetOption

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

func PositionStream(subscriberID string) GetOption

Position allows for getting messages by position subscriber

func SincePosition

func SincePosition(position int64) GetOption

SincePosition allows for getting only more recent messages

func SinceVersion

func SinceVersion(version int64) GetOption

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 MessageHandler interface {
	Type() string
	Process(ctx context.Context, msg Message) error
}

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 Poller

type Poller interface {
	Poll(context.Context) error
}

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

type Subscriber interface {
	Start(context.Context) error
}

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

Directories

Path Synopsis
Package mock_gomessagestore is a generated GoMock package.
Package mock_gomessagestore is a generated GoMock package.
mocks
Package mock_repository is a generated GoMock package.
Package mock_repository is a generated GoMock package.

Jump to

Keyboard shortcuts

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