Observer Pattern #

A new order comes into the e-commerce system. As a result, many things must happen: inventory must be decremented, a confirmation email must be sent, the fulfillment team must be notified, analytics must be updated, and loyalty points must be added to the buyer’s account. Without the Observer Pattern, OrderService has to call all of them explicitly — every new subsystem that needs to know about new orders means one more line of code in OrderService. The OrderService, which should only handle order logic, ends up handling inventory, email, analytics, and loyalty all at once — and every time a new subsystem appears, it has to change. The Observer Pattern reverses this responsibility: OrderService merely announces “a new order has arrived”, and anyone who cares can register to hear that announcement — without OrderService needing to know who is listening or how many.

What Is the Observer Pattern? #

The Observer Pattern is a behavioral design pattern that defines a subscription mechanism to notify many objects about events that occur on the object they are watching, without the watched object (the Subject) needing to know who is watching it or what they do with the information.

Two main roles in the Observer Pattern:

  • Subject (Publisher) — the object that holds state and notifies observers when the state changes; it does not know who its observers are, only that they exist
  • Observer (Subscriber) — the object that registers itself to receive notifications; it is responsible for what happens when a notification arrives

Three properties define the Observer Pattern:

  • Loose coupling — the Subject does not depend on concrete Observers; the two only know each other through an interface
  • Dynamic subscription — Observers can be registered and removed at runtime, without the Subject needing to change
  • Broadcast — one notification from the Subject reaches all registered Observers at once
sequenceDiagram
    participant OS as OrderService (Subject)
    participant IS as InventoryObserver
    participant ES as EmailObserver
    participant AS as AnalyticsObserver

    Note over IS,AS: Register themselves with OrderService
    IS->>OS: Register(inventoryObs)
    ES->>OS: Register(emailObs)
    AS->>OS: Register(analyticsObs)

    Note over OS: A new order arrives
    OS->>OS: CreateOrder(req)
    OS->>IS: Notify(OrderCreatedEvent)
    OS->>ES: Notify(OrderCreatedEvent)
    OS->>AS: Notify(OrderCreatedEvent)

    IS-->>IS: decrement stock
    ES-->>ES: send confirmation email
    AS-->>AS: record event

The Problem It Solves #

The Observer Pattern solves two interrelated problems: tight coupling between the Subject and the parties reacting to its changes, and the Open-Closed Principle being violated every time a new reaction needs to be added.

The Problem: The Subject Knows Too Much #

// ANTI-PATTERN: OrderService calls every subsystem directly
func (s *OrderService) CreateOrder(req CreateOrderRequest) (*Order, error) {
    order, err := s.repo.Create(req)
    if err != nil {
        return nil, err
    }

    // OrderService has to know about inventory
    s.inventorySvc.DeductStock(order.Items)

    // OrderService has to know about email
    s.emailSvc.SendConfirmation(order.UserEmail, order)

    // OrderService has to know about analytics
    s.analyticsSvc.TrackPurchase(order)

    // OrderService has to know about loyalty
    s.loyaltySvc.AddPoints(order.UserID, order.TotalAmount)

    // Add WhatsApp notifications? Modify OrderService again.
    // Add fraud detection? Modify OrderService again.
    // Total coupling: OrderService knows 4+ subsystems

    return order, nil
}

// CORRECT: OrderService only knows that an event exists, not who needs to know
func (s *OrderService) CreateOrder(req CreateOrderRequest) (*Order, error) {
    order, err := s.repo.Create(req)
    if err != nil {
        return nil, err
    }

    // One line — does not care who is listening or how many
    s.eventBus.Publish(OrderCreatedEvent{Order: order})

    return order, nil
}
// Add WhatsApp? Register a WhatsAppObserver — OrderService does not change
// Add fraud detection? Register a FraudObserver — OrderService does not change

Structure and Components #

classDiagram
    class Subject {
        <<interface>>
        +Register(observer Observer)
        +Unregister(observer Observer)
        +Notify(event Event)
    }

    class Observer {
        <<interface>>
        +OnEvent(event Event)
    }

    class OrderEventBus {
        -observers map[EventType][]Observer
        -mu sync.RWMutex
        +Register(eventType EventType, observer Observer)
        +Unregister(eventType EventType, observer Observer)
        +Publish(event Event)
    }

    class InventoryObserver {
        -inventorySvc InventoryService
        +OnEvent(event Event)
    }

    class EmailObserver {
        -emailSvc EmailService
        +OnEvent(event Event)
    }

    class AnalyticsObserver {
        -analyticsSvc AnalyticsService
        +OnEvent(event Event)
    }

    Subject <|.. OrderEventBus
    Observer <|.. InventoryObserver
    Observer <|.. EmailObserver
    Observer <|.. AnalyticsObserver
    OrderEventBus o-- Observer : manages list
ComponentRoleCharacteristics
Subject interfaceContract for registering and notifying observersRegister, Unregister, Notify
Concrete SubjectStores the observer list, triggers notificationsThread-safe with a mutex
Observer interfaceContract every subscriber must satisfyUsually one method: OnEvent or Update
Concrete ObserverReacts to events according to its responsibilityIndependent of each other

Full Implementation: Order Event System #

Event Types and Interfaces #

package events

import "time"

// EventType defines the kinds of events that can be published.
type EventType string

const (
    EventOrderCreated   EventType = "order.created"
    EventOrderPaid      EventType = "order.paid"
    EventOrderShipped   EventType = "order.shipped"
    EventOrderCancelled EventType = "order.cancelled"
)

// Event is the interface for every event in the system.
type Event interface {
    Type() EventType
    OccurredAt() time.Time
}

// OrderCreatedEvent is published when a new order is successfully created.
type OrderCreatedEvent struct {
    OrderID     string
    UserID      string
    UserEmail   string
    Items       []OrderItem
    TotalAmount float64
    CreatedAt   time.Time
}

func (e OrderCreatedEvent) Type() EventType      { return EventOrderCreated }
func (e OrderCreatedEvent) OccurredAt() time.Time { return e.CreatedAt }

// OrderPaidEvent is published when a payment is successfully confirmed.
type OrderPaidEvent struct {
    OrderID       string
    UserID        string
    UserEmail     string
    TransactionID string
    Amount        float64
    PaidAt        time.Time
}

func (e OrderPaidEvent) Type() EventType      { return EventOrderPaid }
func (e OrderPaidEvent) OccurredAt() time.Time { return e.PaidAt }

// OrderShippedEvent is published when an order is shipped.
type OrderShippedEvent struct {
    OrderID        string
    UserID         string
    UserEmail      string
    TrackingNumber string
    Courier        string
    ShippedAt      time.Time
}

func (e OrderShippedEvent) Type() EventType      { return EventOrderShipped }
func (e OrderShippedEvent) OccurredAt() time.Time { return e.ShippedAt }

// OrderItem is an item within an order.
type OrderItem struct {
    ProductID string
    Quantity  int
    Price     float64
}

// Observer is the interface every subscriber must implement.
type Observer interface {
    OnEvent(event Event)
    ObserverName() string // for logging and debugging
}

Concrete Subject: Thread-Safe EventBus #

package events

import (
    "fmt"
    "log/slog"
    "sync"
)

// EventBus is the Subject — it manages the observer list per event type
// and distributes events to the relevant observers.
type EventBus struct {
    mu        sync.RWMutex
    observers map[EventType][]Observer
    logger    *slog.Logger
}

// NewEventBus creates a new, ready-to-use EventBus.
func NewEventBus(logger *slog.Logger) *EventBus {
    return &EventBus{
        observers: make(map[EventType][]Observer),
        logger:    logger,
    }
}

// Register subscribes an observer to a specific event type.
// The same observer can be registered for several event types.
func (b *EventBus) Register(eventType EventType, observer Observer) {
    b.mu.Lock()
    defer b.mu.Unlock()

    // Check whether the observer is already registered for this event type
    for _, existing := range b.observers[eventType] {
        if existing == observer {
            b.logger.Warn("observer already registered",
                "event_type", eventType,
                "observer", observer.ObserverName())
            return
        }
    }

    b.observers[eventType] = append(b.observers[eventType], observer)
    b.logger.Info("observer registered",
        "event_type", eventType,
        "observer", observer.ObserverName(),
        "total", len(b.observers[eventType]),
    )
}

// Unregister removes an observer from a specific event type.
func (b *EventBus) Unregister(eventType EventType, observer Observer) {
    b.mu.Lock()
    defer b.mu.Unlock()

    observers := b.observers[eventType]
    for i, existing := range observers {
        if existing == observer {
            b.observers[eventType] = append(observers[:i], observers[i+1:]...)
            b.logger.Info("observer unregistered",
                "event_type", eventType,
                "observer", observer.ObserverName(),
            )
            return
        }
    }
}

// Publish distributes an event to all registered observers.
// Notification is synchronous — all observers finish before Publish returns.
func (b *EventBus) Publish(event Event) {
    b.mu.RLock()
    observers := make([]Observer, len(b.observers[event.Type()]))
    copy(observers, b.observers[event.Type()])
    b.mu.RUnlock()

    b.logger.Info("publishing event",
        "event_type", event.Type(),
        "observer_count", len(observers),
        "occurred_at", event.OccurredAt(),
    )

    for _, observer := range observers {
        func() {
            defer func() {
                if r := recover(); r != nil {
                    b.logger.Error("observer panicked",
                        "observer", observer.ObserverName(),
                        "event_type", event.Type(),
                        "panic", r,
                    )
                }
            }()
            observer.OnEvent(event)
        }()
    }
}

// PublishAsync distributes an event asynchronously via goroutines.
// Use it for observers that can run in parallel and independently.
func (b *EventBus) PublishAsync(event Event) {
    b.mu.RLock()
    observers := make([]Observer, len(b.observers[event.Type()]))
    copy(observers, b.observers[event.Type()])
    b.mu.RUnlock()

    for _, observer := range observers {
        go func(obs Observer) {
            defer func() {
                if r := recover(); r != nil {
                    b.logger.Error("async observer panicked",
                        "observer", obs.ObserverName(),
                        "event_type", event.Type(),
                        "panic", fmt.Sprintf("%v", r),
                    )
                }
            }()
            obs.OnEvent(event)
        }(observer)
    }
}

// ObserverCount returns the number of observers registered for a specific event type.
func (b *EventBus) ObserverCount(eventType EventType) int {
    b.mu.RLock()
    defer b.mu.RUnlock()
    return len(b.observers[eventType])
}

Concrete Observers #

package observers

import (
    "context"
    "fmt"
    "log/slog"

    "myapp/events"
)

// InventoryObserver decrements stock when an order is created.
type InventoryObserver struct {
    inventorySvc InventoryService
    logger       *slog.Logger
}

func NewInventoryObserver(svc InventoryService, logger *slog.Logger) *InventoryObserver {
    return &InventoryObserver{inventorySvc: svc, logger: logger}
}

func (o *InventoryObserver) OnEvent(event events.Event) {
    e, ok := event.(events.OrderCreatedEvent)
    if !ok {
        return // irrelevant event, ignore it
    }

    ctx := context.Background()
    if err := o.inventorySvc.DeductStock(ctx, e.Items); err != nil {
        o.logger.Error("failed to deduct stock",
            "order_id", e.OrderID,
            "error", err,
        )
        return
    }

    o.logger.Info("stock deducted", "order_id", e.OrderID, "items", len(e.Items))
}

func (o *InventoryObserver) ObserverName() string { return "InventoryObserver" }


// EmailObserver sends confirmation emails when orders are created and paid.
type EmailObserver struct {
    emailSvc EmailService
    logger   *slog.Logger
}

func NewEmailObserver(svc EmailService, logger *slog.Logger) *EmailObserver {
    return &EmailObserver{emailSvc: svc, logger: logger}
}

func (o *EmailObserver) OnEvent(event events.Event) {
    ctx := context.Background()

    switch e := event.(type) {
    case events.OrderCreatedEvent:
        if err := o.emailSvc.SendOrderConfirmation(ctx, e.UserEmail, e.OrderID, e.TotalAmount); err != nil {
            o.logger.Error("failed to send confirmation email",
                "order_id", e.OrderID,
                "email", e.UserEmail,
                "error", err,
            )
        }

    case events.OrderShippedEvent:
        if err := o.emailSvc.SendShippingNotification(ctx, e.UserEmail, e.OrderID, e.TrackingNumber, e.Courier); err != nil {
            o.logger.Error("failed to send shipping email",
                "order_id", e.OrderID,
                "error", err,
            )
        }

    default:
        // Other events are not handled by EmailObserver
    }
}

func (o *EmailObserver) ObserverName() string { return "EmailObserver" }


// AnalyticsObserver records every event to the analytics platform.
type AnalyticsObserver struct {
    analyticsSvc AnalyticsService
    logger       *slog.Logger
}

func NewAnalyticsObserver(svc AnalyticsService, logger *slog.Logger) *AnalyticsObserver {
    return &AnalyticsObserver{analyticsSvc: svc, logger: logger}
}

func (o *AnalyticsObserver) OnEvent(event events.Event) {
    ctx := context.Background()

    properties := map[string]interface{}{
        "event_type":  string(event.Type()),
        "occurred_at": event.OccurredAt(),
    }

    switch e := event.(type) {
    case events.OrderCreatedEvent:
        properties["order_id"] = e.OrderID
        properties["user_id"] = e.UserID
        properties["amount"] = e.TotalAmount
        properties["item_count"] = len(e.Items)

    case events.OrderPaidEvent:
        properties["order_id"] = e.OrderID
        properties["user_id"] = e.UserID
        properties["amount"] = e.Amount
        properties["transaction_id"] = e.TransactionID
    }

    if err := o.analyticsSvc.Track(ctx, string(event.Type()), properties); err != nil {
        o.logger.Warn("failed to track analytics event", "error", err)
        // Non-fatal: analytics failure must not stop the main flow
    }
}

func (o *AnalyticsObserver) ObserverName() string { return "AnalyticsObserver" }


// LoyaltyObserver adds loyalty points when a payment succeeds.
type LoyaltyObserver struct {
    loyaltySvc LoyaltyService
    logger     *slog.Logger
}

func NewLoyaltyObserver(svc LoyaltyService, logger *slog.Logger) *LoyaltyObserver {
    return &LoyaltyObserver{loyaltySvc: svc, logger: logger}
}

func (o *LoyaltyObserver) OnEvent(event events.Event) {
    e, ok := event.(events.OrderPaidEvent)
    if !ok {
        return
    }

    // Calculate points: 1 point per 10,000 paid
    points := int(e.Amount / 10000)
    if points == 0 {
        return
    }

    ctx := context.Background()
    if err := o.loyaltySvc.AddPoints(ctx, e.UserID, points); err != nil {
        o.logger.Error("failed to add loyalty points",
            "user_id", e.UserID,
            "points", points,
            "error", err,
        )
        return
    }

    o.logger.Info("loyalty points added", "user_id", e.UserID, "points", points)
}

func (o *LoyaltyObserver) ObserverName() string { return "LoyaltyObserver" }

Subject: A Clean OrderService #

package order

import (
    "context"
    "fmt"
    "time"

    "myapp/events"
)

// OrderService is the Subject — it publishes events, caring nothing about its observers.
type OrderService struct {
    repo     OrderRepository
    eventBus *events.EventBus
}

func NewOrderService(repo OrderRepository, eventBus *events.EventBus) *OrderService {
    return &OrderService{repo: repo, eventBus: eventBus}
}

func (s *OrderService) CreateOrder(ctx context.Context, req CreateOrderRequest) (*Order, error) {
    order, err := s.repo.Create(ctx, req)
    if err != nil {
        return nil, fmt.Errorf("failed to create order: %w", err)
    }

    // Publish the event — does not know or care who subscribes
    s.eventBus.Publish(events.OrderCreatedEvent{
        OrderID:     order.ID,
        UserID:      req.UserID,
        UserEmail:   req.UserEmail,
        Items:       req.Items,
        TotalAmount: order.TotalAmount,
        CreatedAt:   time.Now(),
    })

    return order, nil
}

func (s *OrderService) ConfirmPayment(ctx context.Context, orderID, transactionID string, amount float64) error {
    order, err := s.repo.UpdateStatus(ctx, orderID, "paid")
    if err != nil {
        return fmt.Errorf("failed to update order status: %w", err)
    }

    s.eventBus.Publish(events.OrderPaidEvent{
        OrderID:       orderID,
        UserID:        order.UserID,
        UserEmail:     order.UserEmail,
        TransactionID: transactionID,
        Amount:        amount,
        PaidAt:        time.Now(),
    })

    return nil
}

func (s *OrderService) ShipOrder(ctx context.Context, orderID, trackingNumber, courier string) error {
    order, err := s.repo.UpdateStatus(ctx, orderID, "shipped")
    if err != nil {
        return fmt.Errorf("failed to update order status: %w", err)
    }

    s.eventBus.Publish(events.OrderShippedEvent{
        OrderID:        orderID,
        UserID:         order.UserID,
        UserEmail:      order.UserEmail,
        TrackingNumber: trackingNumber,
        Courier:        courier,
        ShippedAt:      time.Now(),
    })

    return nil
}

Wiring in main.go #

func main() {
    logger := slog.Default()
    eventBus := events.NewEventBus(logger)

    // Create the observers
    inventoryObs := observers.NewInventoryObserver(inventorySvc, logger)
    emailObs     := observers.NewEmailObserver(emailSvc, logger)
    analyticsObs := observers.NewAnalyticsObserver(analyticsSvc, logger)
    loyaltyObs   := observers.NewLoyaltyObserver(loyaltySvc, logger)

    // Register observers for the relevant events
    eventBus.Register(events.EventOrderCreated, inventoryObs)
    eventBus.Register(events.EventOrderCreated, emailObs)
    eventBus.Register(events.EventOrderCreated, analyticsObs)

    eventBus.Register(events.EventOrderPaid, loyaltyObs)
    eventBus.Register(events.EventOrderPaid, analyticsObs)

    eventBus.Register(events.EventOrderShipped, emailObs)
    eventBus.Register(events.EventOrderShipped, analyticsObs)

    // OrderService does not need to know any of this
    orderSvc := order.NewOrderService(orderRepo, eventBus)

    // Adding a FraudDetectionObserver later:
    // eventBus.Register(events.EventOrderCreated, fraudObs)
    // — no other code changes
}

Async Observer: Non-Blocking Notification #

For observers that can run independently and do not need to wait for each other, async notification via goroutines is more appropriate.

// AsyncEventBus distributes an event to each observer in a separate goroutine.
// Use it when observers do not need a guaranteed execution order.
type AsyncEventBus struct {
    EventBus
    errorHandler func(obs Observer, event Event, err interface{})
}

func NewAsyncEventBus(logger *slog.Logger) *AsyncEventBus {
    bus := &AsyncEventBus{}
    bus.observers = make(map[EventType][]Observer)
    bus.logger = logger
    bus.errorHandler = func(obs Observer, event Event, err interface{}) {
        logger.Error("async observer error",
            "observer", obs.ObserverName(),
            "event_type", event.Type(),
            "error", fmt.Sprintf("%v", err),
        )
    }
    return bus
}

func (b *AsyncEventBus) Publish(event Event) {
    b.mu.RLock()
    observers := make([]Observer, len(b.observers[event.Type()]))
    copy(observers, b.observers[event.Type()])
    b.mu.RUnlock()

    var wg sync.WaitGroup
    for _, observer := range observers {
        wg.Add(1)
        go func(obs Observer) {
            defer wg.Done()
            defer func() {
                if r := recover(); r != nil {
                    b.errorHandler(obs, event, r)
                }
            }()
            obs.OnEvent(event)
        }(observer)
    }
    // Wait for all to finish — or not, depending on the need
    // wg.Wait() // uncomment if you need to wait for all observers to finish
}

Sync vs async comparison:

AspectSynchronous PublishAsynchronous Publish
Execution orderGuaranteed — observers are called in sequenceNot guaranteed — runs in parallel
Error handlingAn observer error can delay the next observerOne observer’s error does not affect others
RollbackRollback can be implemented if an observer failsHard — the event is already published
PerformanceA slow observer blocks the SubjectObservers run in parallel
Best forInterdependent observers that need orderingIndependent, best-effort observers

Preventing Memory Leaks #

One of the biggest mishandled problems of the Observer Pattern is memory leaks — observers that are no longer active remain registered and prevent garbage collection.

// ANTI-PATTERN: observers are never unregistered
type SessionHandler struct {
    eventBus *events.EventBus
}

func (h *SessionHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
    // A new observer is registered on every request — never removed!
    obs := &RequestScopedObserver{requestID: r.Header.Get("X-Request-ID")}
    h.eventBus.Register(events.EventOrderCreated, obs)
    // After the request finishes, obs is not removed → it keeps accumulating
}

// CORRECT: always unregister observers that are no longer needed
func (h *SessionHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
    obs := &RequestScopedObserver{requestID: r.Header.Get("X-Request-ID")}
    h.eventBus.Register(events.EventOrderCreated, obs)
    defer h.eventBus.Unregister(events.EventOrderCreated, obs) // automatic cleanup

    // handle request...
}

// Or: use a WeakReference-style approach with a timeout
type TimedObserver struct {
    observer  Observer
    expiresAt time.Time
}

// EventBus with auto-cleanup of expired observers
func (b *EventBus) CleanupExpired() {
    b.mu.Lock()
    defer b.mu.Unlock()
    now := time.Now()
    for eventType, observers := range b.observers {
        var active []Observer
        for _, obs := range observers {
            if timed, ok := obs.(*TimedObserver); ok {
                if now.Before(timed.expiresAt) {
                    active = append(active, obs)
                }
            } else {
                active = append(active, obs)
            }
        }
        b.observers[eventType] = active
    }
}

Observers Living Longer Than Expected

If the Subject holds a reference to an Observer, the Observer will not be garbage collected even when all other references to it are gone. This is a very common source of memory leaks. Always provide an Unregister mechanism and make sure the Observer lifecycle is managed properly — especially for Observers created per-request or per-session.


Testing the Observer Pattern #

// MockObserver for testing — records every event it receives
type MockObserver struct {
    ReceivedEvents []events.Event
    mu             sync.Mutex
    ShouldPanic    bool
}

func (m *MockObserver) OnEvent(event events.Event) {
    if m.ShouldPanic {
        panic("mock panic")
    }
    m.mu.Lock()
    defer m.mu.Unlock()
    m.ReceivedEvents = append(m.ReceivedEvents, event)
}

func (m *MockObserver) ObserverName() string { return "MockObserver" }

func (m *MockObserver) EventCount() int {
    m.mu.Lock()
    defer m.mu.Unlock()
    return len(m.ReceivedEvents)
}


func TestEventBus_Publish_NotifiesAllObservers(t *testing.T) {
    bus := events.NewEventBus(slog.Default())

    obs1 := &MockObserver{}
    obs2 := &MockObserver{}

    bus.Register(events.EventOrderCreated, obs1)
    bus.Register(events.EventOrderCreated, obs2)

    event := events.OrderCreatedEvent{
        OrderID: "ORD-001", UserID: "user-1",
        TotalAmount: 150000, CreatedAt: time.Now(),
    }
    bus.Publish(event)

    if obs1.EventCount() != 1 {
        t.Errorf("expected obs1 to receive 1 event, got %d", obs1.EventCount())
    }
    if obs2.EventCount() != 1 {
        t.Errorf("expected obs2 to receive 1 event, got %d", obs2.EventCount())
    }
}

func TestEventBus_Unregister_StopsNotification(t *testing.T) {
    bus := events.NewEventBus(slog.Default())
    obs := &MockObserver{}

    bus.Register(events.EventOrderCreated, obs)
    bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-001", CreatedAt: time.Now()})

    bus.Unregister(events.EventOrderCreated, obs)
    bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-002", CreatedAt: time.Now()})

    if obs.EventCount() != 1 {
        t.Errorf("expected 1 event after unregister, got %d", obs.EventCount())
    }
}

func TestEventBus_PanicInObserver_DoesNotStopOthers(t *testing.T) {
    bus := events.NewEventBus(slog.Default())

    panicObs := &MockObserver{ShouldPanic: true}
    normalObs := &MockObserver{}

    bus.Register(events.EventOrderCreated, panicObs)
    bus.Register(events.EventOrderCreated, normalObs)

    // Publish must not panic even when one observer panics
    require.NotPanics(t, func() {
        bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-001", CreatedAt: time.Now()})
    })

    // The normal observer still receives the event
    if normalObs.EventCount() != 1 {
        t.Errorf("normal observer should still receive event despite panic in another observer")
    }
}

func TestEventBus_TypeFiltering(t *testing.T) {
    bus := events.NewEventBus(slog.Default())
    obs := &MockObserver{}

    // Only subscribe to EventOrderPaid
    bus.Register(events.EventOrderPaid, obs)

    // Publish EventOrderCreated — obs should not receive it
    bus.Publish(events.OrderCreatedEvent{OrderID: "ORD-001", CreatedAt: time.Now()})
    if obs.EventCount() != 0 {
        t.Error("observer should not receive event for unsubscribed type")
    }

    // Publish EventOrderPaid — should be received
    bus.Publish(events.OrderPaidEvent{OrderID: "ORD-001", Amount: 100000, PaidAt: time.Now()})
    if obs.EventCount() != 1 {
        t.Errorf("expected 1 event for subscribed type, got %d", obs.EventCount())
    }
}

func TestEventBus_ConcurrentPublish(t *testing.T) {
    bus := events.NewEventBus(slog.Default())
    obs := &MockObserver{}
    bus.Register(events.EventOrderCreated, obs)

    const goroutines = 100
    var wg sync.WaitGroup
    for i := 0; i < goroutines; i++ {
        wg.Add(1)
        go func(i int) {
            defer wg.Done()
            bus.Publish(events.OrderCreatedEvent{
                OrderID:   fmt.Sprintf("ORD-%d", i),
                CreatedAt: time.Now(),
            })
        }(i)
    }
    wg.Wait()

    if obs.EventCount() != goroutines {
        t.Errorf("expected %d events from concurrent publish, got %d", goroutines, obs.EventCount())
    }
}

Observer vs Event Bus vs Mediator #

The Observer Pattern can be implemented in several ways. It is important to understand the differences.

AspectClassic ObserverEvent Bus (Typed)Mediator
Does the Subject know its Observers?Yes — it stores a listNo — only knows EventTypeNo — everything goes through the mediator
Do Observers know the Subject?No — only know the eventNo — only know EventTypeNo
CouplingSubject-Observer loosely coupledMost loosely coupledComponents do not know each other
Event filteringObservers filter themselvesFilter by type at the busThe Mediator decides who needs to know
Best forOne subject, several observersMany subjects, many observersComplex communication between components

When to Use and When Not to #

USE Observer if:
  ✓ A change in one object needs to be known by other objects without tight coupling
  ✓ The number and identity of observers are not known at compile time
  ✓ Observers can be registered and removed dynamically
  ✓ You are building an event-driven or reactive system
  ✓ You want to apply the Open-Closed Principle to reactions to events

AVOID Observer if:
  ✗ Notification chains get too long and make debugging difficult
  ✗ Observers need to know about other observers — use a Mediator
  ✗ Observer order is critical and cannot be guaranteed
  ✗ The system is very small with 1-2 reactions that will never change

Observer Review Checklist #

DESIGN:
  □ The Subject does not depend on concrete Observers — only on the interface
  □ Observers can be registered and removed without changing the Subject
  □ Events carry enough information that Observers do not need to query the Subject back
  □ Each Observer handles only relevant events — ignores the rest

IMPLEMENTATION:
  □ The EventBus is thread-safe — Register, Unregister, Publish use a mutex
  □ Publish copies the observer list before calling (avoid deadlock)
  □ A panic in one observer does not stop notifications to other observers
  □ Observer errors are not propagated to the Subject — the Subject does not care

LIFECYCLE:
  □ There is an easy-to-use Unregister mechanism
  □ Observers are unregistered when no longer needed
  □ No circular references between Subject and Observer

TESTING:
  □ A MockObserver records all events for assertions
  □ Test that all observers are notified
  □ Test that Unregister stops notifications
  □ Test that a panic in one observer does not affect others
  □ Test thread-safety with concurrent Publish

Summary #

  • Observer defines a subscription mechanism — the Subject announces changes, Observers react; the two do not need to know each other concretely.
  • True loose coupling: the Subject does not know who its observers are, how many, or what they do — it only calls OnEvent() through an interface.
  • A typed Event Bus is an evolution of the classic Observer — events are filtered by type at the bus, not in each observer; more scalable for systems with many events.
  • Thread safety is mandatory: Register, Unregister, and Publish are often called from different goroutines; use sync.RWMutex and copy the observer list before iterating.
  • Panic recovery inside Publish — one panicking observer must not stop notifications to other observers; wrap every observer call with recover().
  • Memory leaks are a real risk — unregistered observers prevent garbage collection; always use defer Unregister() for observers with a limited lifecycle.
  • Sync vs Async: synchronous for observers that need ordering; async (goroutines) for independent observers that can run in parallel.
  • Distinguish it from Mediator: Observer is for “one subject, many subscribers”; Mediator is for “many components communicating through an intermediary”.

← Previous: Strategy   Next: Command →

About | Author | Content Scope | Editorial Policy | Privacy Policy | Disclaimer | Contact