Documentation
¶
Overview ¶
Package external_sink integrates external event sources into the eventsourced framework.
When events originate from external systems (e.g., other services publishing to AMQP), the sink consumes them, wraps each event in an eventsourced.ExternalEvent, and processes them through the standard command handler flow. This stores external events in the same event store as internal events, enabling unified event streams.
Architecture ¶
External events flow through these steps:
- External service publishes event to AMQP
- Sink's AMQP consumer receives the event
- Event type is matched to a registered CommandFunc
- CommandFunc creates an ExternalCommand wrapping the event
- Command handler processes the command, producing an eventsourced.ExternalEvent
- ExternalEvent is stored in the event store and published to the internal stream
All external events are stored under a single ExternalAggregate with the fixed ID "external".
Usage ¶
sink, err := external_sink.New("my-service", eventStore,
map[string]any{"User.Created": &ExternalUserCreated{}},
)
if err != nil {
log.Fatal(err)
}
// Register AMQP setup with the connection
err = conn.Start(ctx, append(setups, sink.Setup(goamqp.WithDeadLetter())...)...)
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type CommandFunc ¶
type CommandFunc[T eventsourced.Event] func(source string, event T) ExternalCommand[T]
CommandFunc is a function that maps an external event from a given source to an ExternalCommand. It is called when an external event is received from AMQP.
type ExternalAggregate ¶
type ExternalAggregate struct {
// contains filtered or unexported fields
}
ExternalAggregate is a special aggregate used to track external event processing. All external events are stored under a single aggregate with the fixed ID "external".
func (*ExternalAggregate) Apply ¶
func (e *ExternalAggregate) Apply(_ eventsourced.Event) error
func (*ExternalAggregate) Identity ¶
func (e *ExternalAggregate) Identity() *eventsourced.ID
func (*ExternalAggregate) SetIdentity ¶
func (e *ExternalAggregate) SetIdentity(id eventsourced.ID)
type ExternalCommand ¶
type ExternalCommand[T eventsourced.Event] struct { Source string Type string ExternalEvent *T }
ExternalCommand wraps an external event as a command that produces an eventsourced.ExternalEvent when processed by the command handler.
func (ExternalCommand[T]) Event ¶
func (e ExternalCommand[T]) Event(_ context.Context) eventsourced.Event
func (ExternalCommand[T]) Validate ¶
func (e ExternalCommand[T]) Validate(_ context.Context, _ eventsourced.Aggregate) error
type Sink ¶
type Sink struct {
// contains filtered or unexported fields
}
Sink integrates external event sources into the eventsourced framework. It consumes events from an AMQP stream, wraps them in ExternalEvent, and processes them through the eventsourced command handler.
A single CommandHandler is created lazily on the first message and reused for all subsequent messages, so the ExternalAggregate event history is loaded once instead of replayed per message. The handler is also serialized behind a mutex because commandHandler keeps mutable state (lastSequenceNo, isDeleted, etc.) that is not safe to share across concurrent Handle calls.
func New ¶
func New(serviceName string, eventStore eventsourced.EventStore, types map[string]any, opts ...eventsourced.Option) (*Sink, error)
New creates an external event sink for the given service. types maps the routing key each external event arrives on to its event type, passed as a pointer implementing eventsourced.Event and eventsourced.ExternalSource. The same key is used when the event is re-published on the internal stream, so consumers of StreamName() bind the same domain keys as consumers of the events exchange. Each type is validated for correct JSON struct tags at creation time.
func (*Sink) Setup ¶
func (s *Sink) Setup(opts ...goamqp.ConsumerOptions) []goamqp.Setup
Setup returns the goamqp setup functions for configuring the AMQP stream publisher and the event stream consumer. The consumer's queue is bound to each routing key registered with New, so the broker only delivers the events the sink handles. opts are applied to the consumer, e.g. goamqp.WithDeadLetter() to keep deliveries that fail terminally (such as unparsable payloads) instead of dropping them.
func (*Sink) StreamName ¶
StreamName returns the AMQP stream name used by this sink. The stream name is derived from the service name: "{serviceName}.internal".
func (*Sink) TypeMapper ¶ added in v0.4.0
func (s *Sink) TypeMapper() goamqp.TypeMapper
TypeMapper resolves the routing keys registered with New to their event types. Events re-published on StreamName() carry the same keys, so consumers of that stream can decode with it.