Skip to documentation
Documentation navigation

Documentation navigation

Documentation / guides

Subscriptions

Consume Session events without confusing delivery with persistence.

developer

Session.SubscribeEvents attaches a consumer to the session hub. The filter is evaluated before the subscription’s bounded egress send. Delivery is a live view, not proof that a value was persisted.

Filter contract

The exact filter types are:

type EventFilter struct {
	Ephemeral LoopScope
	Enduring  LoopScope
}

type LoopScope struct {
	All   bool
	Loops map[uuid.UUID]struct{}
}

func (s LoopScope) Matches(loopID uuid.UUID) bool
func ShouldDeliver(filter EventFilter, ev Event) bool

Public session-scoped events bypass loop selection. Public loop-scoped events select Ephemeral or Enduring by the event’s class. Internal events never enter ordinary subscriptions. A zero LoopScope matches no loop; use All: true or a non-empty UUID set deliberately.

Event classDelivery.JournalSeqOverflow policy
Ephemeral0drop for this slow subscriber
Enduringassigned ledger sequencefail the subscriber with SubscriptionLossError

Delivery contract

The exact consumer interface is:

type Subscription interface {
	Events() <-chan Delivery
	Close() error
	Err() error
}

type Delivery struct {
	Event      Event
	JournalSeq uint64
}

The concrete hub channel is bounded to 256 values. Close is idempotent and intentional, so Err remains nil. If an Enduring value cannot enter the buffer, the hub closes the channel and stores *hub.SubscriptionLossError in Err. A consumer must resync from durable history after that error.

Bounded lifecycle

%%{init: {"theme":"dark"}}%%
flowchart LR
    P[Hub publication] --> V{Public event?}
    V -- no --> X[not delivered]
    V -- yes --> F{Filter matches?}
    F -- no --> X
    F -- yes --> B{256-slot buffer available?}
    B -- yes --> D[Events channel]
    B -- no, Ephemeral --> E[drop this delivery]
    B -- no, Enduring --> L[close with SubscriptionLossError]
    D --> C[consumer calls Close]

SessionStopped is an event, not a subscription terminator. The channel closes only on intentional close, hub-forced loss, or hub teardown. A subscriber that sees a closed channel must inspect Err before declaring a clean end.

Consumer example

func consume(ctx context.Context, s session.Session, loopID uuid.UUID) error {
	sub, err := s.SubscribeEvents(event.EventFilter{
		Ephemeral: event.LoopScope{Loops: map[uuid.UUID]struct{}{loopID: {}}},
		Enduring:  event.LoopScope{All: true},
	})
	if err != nil {
		return err
	}
	defer sub.Close()

	for {
		select {
		case <-ctx.Done():
			return ctx.Err()
		case d, ok := <-sub.Events():
			if !ok {
				return sub.Err()
			}
			if d.Event.Class() == event.Enduring {
				log.Printf("durable seq=%d type=%T", d.JournalSeq, d.Event)
			}
		}
	}
}

The last handled Enduring journal sequence is the correct replay cursor. It is distinct from Header.EventID, which remains the event’s idempotency and causal identity.

Source and proof

← back to documentation