Docs / ArchitectureKit / Guides / Publishing Events

Publishing Events

Other systems often need to learn about events, for example to send an email once a book has been borrowed, to call a webhook, to fill a search index, or to tell another team. How they learn about them depends on whether they can read the database themselves.

Observing the Database#

If the other system can read EventSourcingDB, let it observe the events itself. It then gets every event in the order it was stored, and after a restart, it resumes from the last event it has seen. Nothing needs to be forwarded (see Observing Events).

Publishing with a Projection#

If it can not, because it is a mail server, a webhook, or a message broker, or if it should not depend on how the events of your application look, forward the events with a projection that publishes them instead of writing a view. Do not publish from where a command is executed: if the application stops right after Execute, the event is stored, but never published, and nothing tells.

A projection that publishes needs to be resumable (see Resuming Projections), and it reports a failure of the other system as an error of the category ErrTransient:

type LoanMailer struct {
  *architecturekit.TypedProjection
  checkpoints CheckpointStore
}

func NewLoanMailer(mailer Mailer, checkpoints CheckpointStore) *LoanMailer {
  return &LoanMailer{
    TypedProjection: architecturekit.NewProjection().
      On(func(ctx context.Context, event architecturekit.Envelope[BookBorrowed]) error {
        // The ID of the event lets the receiver recognize a mail it got before.
        err := mailer.Send(ctx, event.ID, event.Data.BorrowedBy, "You have borrowed a book.")
        if err != nil {
          return fmt.Errorf("%w: sending a mail: %v", architecturekit.ErrTransient, err)
        }
        return nil
      }),
    checkpoints: checkpoints,
  }
}

func (m *LoanMailer) Checkpoint(ctx context.Context) (string, error) {
  return m.checkpoints.Load(ctx, "loan-mailer")
}

func (m *LoanMailer) SaveCheckpoint(ctx context.Context, eventID string) error {
  return m.checkpoints.Save(ctx, "loan-mailer", eventID)
}

Here, Mailer and CheckpointStore stand for whatever your application uses to send mails and to keep a value. Start the projection as any other (see Starting Projections):

run := architecturekit.StartProjection(ctx, store, "/books", true, NewLoanMailer(mailer, checkpoints))

How a Publisher Behaves#

Since publishing rides on a projection, it behaves like one:

  • It resumes only with a checkpoint. A projection that is not resumable starts from the first event whenever the application starts, which suits a view held in memory. A publisher would send every event again on every start. Within a run, the position is kept either way, so reconnecting repeats nothing.
  • An event may be published twice. The checkpoint is saved after an event has been applied, so an event published right before the application stopped is published again after the restart. Hand over the ID of the event, so that the receiver can recognize an event it got before.
  • A transient failure is tried again. For an error of the category ErrTransient, the run tries the failed event again, with a growing delay. Any other error ends the run.
  • Events arrive in order, one at a time, as they were stored.

Note that the mode of a projection depends on the functions it implements. To make sure that a publisher is resumable, check its mode in a test (see Testing Projections).