broadcasting

package
v0.33.0 Latest Latest
Warning

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

Go to latest
Published: Sep 8, 2026 License: MIT Imports: 10 Imported by: 0

Documentation

Overview

Package broadcasting publishes an event to the clients listening on a channel, and authorizes the subscriptions that ask to listen.

An event is put on its way by BroadcastManager.Queue, which pushes it as a job unless the event asked to go now. The job resolves a driver by name and hands it the channels, the event name and the payload; the drivers live in github.com/arandu-io/hesape/broadcasting/broadcasters.

A subscription arrives at BroadcastController.Authenticate, which asks the driver to authorize it. The decision is made by a Policy and comes back as an auth.Grant, and the name the client is finally signed onto is built from that Grant: every published channel is "<tenant>:<channel>", and the tenant comes off the Grant and from nowhere else. See TenantChannel.

Index

Constants

View Source
const (
	// AuthRoute authorizes a subscription to one channel.
	AuthRoute = "/broadcasting/auth"
	// UserAuthRoute authenticates the connection itself.
	UserAuthRoute = "/broadcasting/user-auth"
)

The two paths BroadcastManager.Routes and BroadcastManager.UserRoutes register.

View Source
const (
	// PrivateChannelPrefix marks a channel that has to be authorized.
	PrivateChannelPrefix = "private-"
	// PresenceChannelPrefix marks a channel that has to be authorized and
	// reports who is listening.
	PresenceChannelPrefix = "presence-"
	// EncryptedPrivateChannelPrefix marks a private channel whose payloads are
	// encrypted end to end.
	EncryptedPrivateChannelPrefix = "private-encrypted-"
)

The prefixes the three kinds of guarded channel carry in front of a name.

They are written by the constructors below and read back off an incoming channel name by UsePusherChannelConventions. They are named here so the writing side and the reading side cannot drift apart.

View Source
const ChannelJoin auth.Action = "broadcasting.join"

ChannelJoin is the auth.Action every channel authorization is issued for.

A decision is made by a Policy and produces an auth.Grant, and a Grant is issued for an action -- so listening on a channel is an action with a name, and auth.Grant.Check refuses a Grant issued for anything else.

It lives here rather than beside the abstract broadcaster because BroadcastController.Authenticate has to check the Grant a driver answered, and this package cannot import the drivers -- they import it.

View Source
const ChannelNameField = "channel_name"

ChannelNameField is the field the socket client sends the channel it wants to listen on under.

View Source
const SocketIDHeader = "X-Socket-ID"

SocketIDHeader is the header BroadcastManager.Socket reads the socket id of the calling connection off.

View Source
const TenantSeparator = ":"

TenantSeparator is the ':' between the tenant and the channel in a published channel name, and the only character a channel name may not contain.

It is named because both halves of TenantChannel and CutTenant spell it, and a separator spelled twice is a separator that drifts.

Variables

View Source
var ErrNoTenant = errors.New("broadcasting: the Grant carries no tenant, and a channel without one belongs to every customer")

ErrNoTenant is returned when the Grant carries no tenant, or carries one that cannot be part of a channel name.

A channel published without a tenant in its name is a channel every customer of the system can subscribe to.

View Source
var ErrTenantInChannelName = errors.New("broadcasting: a channel name may not contain '" + TenantSeparator + "', which is where the tenant goes")

ErrTenantInChannelName is returned when a channel name carries TenantSeparator, which means somebody is naming a tenant.

A tenant is added once, by TenantChannel, out of the Grant. A name that already contains ':' either came from a client choosing whose events it hears, or from a publisher prefixing a second time; both are refused here rather than concatenated.

Without the refusal, an application has to register the pattern "{tenant}:private-orders.{orderId}" for its channels to match at all, and Broadcaster.ExtractAuthParameters then hands the handler a tenant taken out of the request rather than off the Grant.

Functions

func CutTenant

func CutTenant(name string) (tenant, channel string, found bool)

CutTenant is the inverse of TenantChannel: it splits the name that goes on the wire back into the tenant and the channel.

found is false when there is no tenant in front, and then channel is the name unchanged -- which is the shape of every name a client sends, because a client never names a tenant.

The left half has to be a tenant auth.ValidTenant accepts, so a channel that merely contains a colon is not read as one. TenantChannel refuses to build such a name, and this is the reading side of the same rule.

func TenantChannel

func TenantChannel(g auth.Grant, c Channel) (string, error)

TenantChannel is the name a channel is actually published under: "<tenant>:<name>".

The tenant comes from the Grant, never from the path, the body, the query or a header, and every key an application writes -- a cache key, a storage path, a scheduler lock -- is prefixed by it. A channel is the same kind of key with a subscriber on the other end: without the prefix, "private-orders.17" is one channel shared by every customer who has an order 17, and the first one to subscribe reads the others' events.

Both sides of the wire go through this function: the publisher builds the name it publishes on, and the authentication endpoint answers about the same name, so a client that asks for "private-orders.17" is authorized for "acme:private-orders.17" and never gets to choose the "acme".

The zero Grant reaches no channel, which is the answer auth.Grant.Check gives for the same reason: a caller who authorized nothing has no tenant to build a name out of.

func TenantChannels

func TenantChannels(g auth.Grant, channels []Channel) ([]string, error)

TenantChannels is TenantChannel over a list, which is what every driver's Broadcast does first.

Types

type AnonymousEvent

type AnonymousEvent struct {
	InteractsWithBroadcasting
	InteractsWithSockets
	// contains filtered or unexported fields
}

AnonymousEvent is a broadcast with no event type behind it.

It is built by BroadcastManager.On, BroadcastManager.Private and BroadcastManager.Presence, because the constructor is the only place the channels are set and the dispatcher it sends through comes from the manager.

manager.Private("orders.17").As("OrderShipped").With(map[string]any{"total": 42}).Send()

func NewAnonymousEvent

func NewAnonymousEvent(events Dispatcher, channels ...Channel) *AnonymousEvent

NewAnonymousEvent builds an anonymous broadcast over the channels it goes out on. The dispatcher is the one AnonymousEvent.Send hands the broadcast to.

func (*AnonymousEvent) As

func (e *AnonymousEvent) As(name string) *AnonymousEvent

As is the name the event should be broadcast as.

func (*AnonymousEvent) BroadcastAs

func (e *AnonymousEvent) BroadcastAs() string

BroadcastAs is the name given to AnonymousEvent.As, or "AnonymousEvent" when none was.

func (*AnonymousEvent) BroadcastOn

func (e *AnonymousEvent) BroadcastOn() []Channel

BroadcastOn is the channels the event goes out on.

func (*AnonymousEvent) BroadcastWith

func (e *AnonymousEvent) BroadcastWith() map[string]any

BroadcastWith is the payload given to AnonymousEvent.With, or an empty map.

func (*AnonymousEvent) Send

func (e *AnonymousEvent) Send() []any

Send broadcasts the event: it builds the PendingBroadcast, applies the connection and the socket exclusion, and hands it to the dispatcher.

func (*AnonymousEvent) SendNow

func (e *AnonymousEvent) SendNow() []any

SendNow broadcasts the event synchronously.

It sets the flag AnonymousEvent.ShouldBroadcastNow answers and sends; BroadcastManager.Queue is what reads the flag and skips the queue.

func (*AnonymousEvent) ShouldBroadcastNow

func (e *AnonymousEvent) ShouldBroadcastNow() bool

ShouldBroadcastNow is true when the event was sent with AnonymousEvent.SendNow.

func (*AnonymousEvent) ToOthers

func (e *AnonymousEvent) ToOthers(socket string) *AnonymousEvent

ToOthers broadcasts to everyone except the current user.

The socket id is the argument and is kept until Send, because there is no ambient request to read it off -- see InteractsWithSockets.

func (*AnonymousEvent) Via

func (e *AnonymousEvent) Via(connection string) *AnonymousEvent

Via is the connection the event should be broadcast on.

It is not InteractsWithBroadcasting.BroadcastVia: this one only records the name, and Send is what hands it on. Keeping the two apart is what makes ToOthers and Via composable in either order.

func (*AnonymousEvent) With

func (e *AnonymousEvent) With(payload any) *AnonymousEvent

With is the payload the event should be broadcast with.

Go has no union type, so it takes any and reads two shapes: an Arrayable is flattened with ToArray, and a map[string]any has each of its values flattened the same way. Anything else is ignored.

type Arrayable

type Arrayable interface {
	// ToArray flattens the value into what it travels as.
	ToArray() map[string]any
}

Arrayable is the one thing a payload value is checked for before it travels: a value that knows how to flatten itself.

type BroadcastController

type BroadcastController struct {
	// contains filtered or unexported fields
}

BroadcastController is the two endpoints a socket client calls before it is allowed to listen.

func NewBroadcastController

func NewBroadcastController(broadcast *BroadcastManager) *BroadcastController

NewBroadcastController builds the controller over a manager.

func (*BroadcastController) Authenticate

func (c *BroadcastController) Authenticate(w http.ResponseWriter, r *http.Request)

Authenticate authorizes the request for channel access.

This is where a private channel becomes an authorization decision. The channel name arrives from the client, the subject arrives from the context where the session middleware put it, and the driver's Auth runs the channel's Policy through auth.Authorize -- so a refusal is auth.ErrForbidden and a success is a Grant, exactly as it is on the way into a repository. The client never names the tenant: it comes off the Grant.

The refusal is 403, and the body is deliberately the same sentence for every refusal: the reason a channel was denied is a fact about somebody else's data.

The Grant is checked for ChannelJoin before anything is written. A driver that decides nothing answers the zero Grant with no error at all -- which is what the log and null drivers do -- and the check is what keeps that from reading as a success.

func (*BroadcastController) AuthenticateUser

func (c *BroadcastController) AuthenticateUser(w http.ResponseWriter, r *http.Request)

AuthenticateUser authenticates the current user for the connection itself, rather than for one channel.

See https://pusher.com/docs/channels/server_api/authenticating-users for the document the client expects.

A driver that resolves no user is a 403. So is a driver that has no resolver at all: the endpoint exists to answer who somebody is, and a deployment that never registered a way to say cannot answer.

type BroadcastError

type BroadcastError struct {
	// contains filtered or unexported fields
}

BroadcastError is what a driver returns when a broadcast could not be published.

The cause is kept rather than flattened into the message, because errors.Is on the driver's error is how a caller tells "the broker is down" from "the payload would not encode".

func NewBroadcastError

func NewBroadcastError(format string, args ...any) *BroadcastError

NewBroadcastError builds the error with a formatted message.

func WrapBroadcastError

func WrapBroadcastError(cause error, format string, args ...any) *BroadcastError

WrapBroadcastError builds the error over the driver failure that caused it. The message is formatted as NewBroadcastError formats it.

func (*BroadcastError) Error

func (e *BroadcastError) Error() string

Error implements error with the message the error was built with.

func (*BroadcastError) Unwrap

func (e *BroadcastError) Unwrap() error

Unwrap exposes the driver failure underneath to errors.Is and errors.As. It is nil when the error was raised without one.

type BroadcastEvent

type BroadcastEvent struct {
	// Event is what is being broadcast.
	Event any
	// Tries is how many times the job may be attempted.
	Tries int
	// Timeout is how long one attempt may run. It is a Duration rather than a
	// count of seconds, because an API that measures time in bare ints is one
	// that gets milliseconds passed to it.
	Timeout time.Duration
	// Backoff is how long to wait before retrying.
	Backoff time.Duration
	// MaxExceptions is how many uncaught failures the job may accumulate before
	// it is given up on.
	MaxExceptions int
	// DeleteWhenMissingModels tells the queue to drop the job when a model it
	// carries can no longer be found. It starts true.
	DeleteWhenMissingModels bool
}

BroadcastEvent is the queued job that carries an event to the broadcasters.

The fields are exported because the queue reads them off the job. The constructor fills each one from the optional interfaces below, and an application that wants a different value sets the field.

func NewBroadcastEvent

func NewBroadcastEvent(event any) *BroadcastEvent

NewBroadcastEvent builds the job that carries event to the broadcasters.

func (*BroadcastEvent) Clone

func (b *BroadcastEvent) Clone() *BroadcastEvent

Clone copies the job and the event it carries, so a queued copy cannot be mutated by whoever still holds the original.

The copy is shallow: a pointer event is followed one level and the struct behind it copied.

func (*BroadcastEvent) DisplayName

func (b *BroadcastEvent) DisplayName() string

DisplayName names the event being carried, which is what a worker log line has to carry for anyone to know which event failed.

It is reflect.Type.String of the event, and the empty string when there is none.

func (*BroadcastEvent) Failed

func (b *BroadcastEvent) Failed(ctx context.Context, cause error) error

Failed hands the failure to the event, when the event wants it.

It returns an error because an event's own failure handling can fail, and swallowing that would lose the only report of it.

func (*BroadcastEvent) Handle

func (b *BroadcastEvent) Handle(ctx context.Context, g auth.Grant, manager Factory) error

Handle is the queued job's body: it names the event, reads its payload and publishes it on every connection the event asked for.

ctx is there because publishing to a broker is I/O.

g is where the tenant comes from. Every channel this job publishes on is named "<tenant>:<channel>", the tenant comes from the Grant and from nothing else, and a job that could publish without one would publish into a channel every customer of the system can subscribe to. The Grant is the job's own -- in a worker, queue/jobs.GrantFor rebuilds exactly the Grant the push authorized.

An event on no channels returns without touching a driver.

func (*BroadcastEvent) Middleware

func (b *BroadcastEvent) Middleware() []any

Middleware is the job middleware of the underlying event, or none when it declares none.

type BroadcastManager

type BroadcastManager struct {
	// contains filtered or unexported fields
}

BroadcastManager resolves drivers by name, caches them, and is the entry point for everything that leaves this package.

A BroadcastManager is safe for concurrent use.

func NewBroadcastManager

func NewBroadcastManager(config Config, events Dispatcher, jobs Queue, locks UniqueLock) *BroadcastManager

NewBroadcastManager builds a manager over everything it needs: the configuration, the dispatcher an anonymous broadcast is sent through, the queue and the lock.

The queue and the lock may be nil: an application that only ever broadcasts synchronously never reaches them, and one that does is told so by name rather than by a nil dereference.

func (*BroadcastManager) ChannelRoutes

func (m *BroadcastManager) ChannelRoutes(r Router)

ChannelRoutes is an alias of BroadcastManager.Routes.

func (*BroadcastManager) Connection

func (m *BroadcastManager) Connection(name string) (Broadcaster, error)

Connection is BroadcastManager.Driver under the name the Factory contract gives it.

func (*BroadcastManager) Driver

func (m *BroadcastManager) Driver(name string) (Broadcaster, error)

Driver returns a driver instance, resolving and caching it the first time. An empty name is the default driver.

func (*BroadcastManager) Event

func (m *BroadcastManager) Event(event any) *PendingBroadcast

Event begins broadcasting an event.

The returned broadcast reaches the dispatcher through PendingBroadcast.Send, which has to be called: nothing sends it on the way out of scope.

func (*BroadcastManager) Extend

func (m *BroadcastManager) Extend(driver string, creator DriverCreator) *BroadcastManager

Extend registers a driver creator under a name.

It is also how the three drivers this ecosystem carries are registered -- see [BroadcastManager.resolve] for why they are not methods on this type.

func (*BroadcastManager) ForgetDrivers

func (m *BroadcastManager) ForgetDrivers() *BroadcastManager

ForgetDrivers forgets every resolved driver.

func (*BroadcastManager) GetDefaultDriver

func (m *BroadcastManager) GetDefaultDriver() string

GetDefaultDriver is the connection name used when none is given.

func (*BroadcastManager) On

func (m *BroadcastManager) On(channels ...Channel) *AnonymousEvent

On begins an anonymous broadcast to the given channels.

func (*BroadcastManager) Presence

func (m *BroadcastManager) Presence(channel string) *AnonymousEvent

Presence begins an anonymous broadcast to a presence channel.

func (*BroadcastManager) Private

func (m *BroadcastManager) Private(channel string) *AnonymousEvent

Private begins an anonymous broadcast to a private channel.

func (*BroadcastManager) Purge

func (m *BroadcastManager) Purge(name string)

Purge forgets one resolved driver, so the next call to Driver builds it again. An empty name is the default driver.

func (*BroadcastManager) Queue

func (m *BroadcastManager) Queue(ctx context.Context, g auth.Grant, event any) error

Queue puts the event on its way.

An event that asked to go now is handled in this goroutine, through BroadcastEvent.Handle. Everything else is pushed, and an event that asked to be unique takes a lock first and is dropped when somebody already holds it. The event is cloned either way: the copy that travels must not change under whoever still holds the original.

ctx and g are what BroadcastEvent.Handle documents: the I/O, and the Grant the tenant comes off.

func (*BroadcastManager) Routes

func (m *BroadcastManager) Routes(r Router)

Routes registers the channel authentication endpoint.

Middleware is attached by whoever owns the router, and so is the CSRF exemption this endpoint needs: it is reached by the socket client with its own credentials, and a CSRF token check on it rejects every legitimate subscription.

func (*BroadcastManager) SetDefaultDriver

func (m *BroadcastManager) SetDefaultDriver(name string)

SetDefaultDriver sets the connection name used when none is given.

func (*BroadcastManager) Socket

func (m *BroadcastManager) Socket(r *http.Request) string

Socket is the socket id of the connection that sent this request, read off the X-Socket-ID header.

The request is an argument because Go has no ambient one, and a nil request answers the empty string.

func (*BroadcastManager) UserRoutes

func (m *BroadcastManager) UserRoutes(r Router)

UserRoutes registers the user authentication endpoint.

type Broadcaster

type Broadcaster interface {
	// Auth decides whether the caller may listen on a channel, and answers the
	// body the client is sent back.
	//
	// The Grant comes back beside the response, and that is the whole point of
	// this framework: a private channel is an authorization decision, the
	// decision is made by a Policy, and the Grant is the proof it happened.
	// auth.ErrForbidden is the refusal, and nothing scoped by tenant is
	// reachable without the Grant -- including the name the event is finally
	// published under, which is [TenantChannel](grant, channel).
	//
	// The subject is not an argument. It comes from the context, where the
	// middleware that loaded the session put it -- auth.SubjectFrom(ctx) --
	// which, unlike a parameter, cannot be supplied by the caller being
	// authorized.
	//
	// The channel is the raw name the client asked for, prefix and all:
	// "private-orders.17". The driver normalizes it, and refuses it outright if
	// it names a tenant -- see [RequestedChannel].
	//
	// The Grant is the answer, not the response. A driver that authorized
	// nobody answers the zero Grant, and the zero Grant fails
	// auth.Grant.Check([ChannelJoin]); a caller that reads only the response is
	// reading a body a driver may have produced without deciding anything.
	Auth(ctx context.Context, channel string) (auth.Grant, any, error)

	// ValidAuthenticationResponse is the body a successful Auth answers with,
	// built from whatever the channel handler returned.
	//
	// It takes the Grant rather than the request, because everything it needs --
	// who the subject is, which tenant they are in -- is on the Grant, and
	// reading it there is the difference between a response about the person who
	// was authorized and a response about whoever sent the bytes.
	//
	// channel is the channel the client asked for, without a tenant. It is a
	// parameter because the answer has to say what it is an answer about: a
	// relay given a bare `true` signs the socket onto the string the client
	// sent. The implementations name the channel out of [TenantChannel](g,
	// channel) and put that name in the body, which is why this method cannot be
	// called without a Grant.
	ValidAuthenticationResponse(ctx context.Context, g auth.Grant, channel Channel, result any) (any, error)

	// Broadcast publishes the event.
	//
	// ctx is there because this is I/O. Every channel goes out under
	// [TenantChannel](g, channel), so the tenant is in the name of everything
	// published and a subscriber in one tenant cannot name a channel in another.
	Broadcast(ctx context.Context, g auth.Grant, channels []Channel, event string, payload map[string]any) error
}

Broadcaster is the three methods every driver has.

The implementations live in github.com/arandu-io/hesape/broadcasting/broadcasters. The contract stays here because the package that implements it imports this one.

type BroadcastsAs

type BroadcastsAs interface {
	BroadcastAs() string
}

BroadcastsAs is the name the event goes out under. The fallback is BroadcastEvent.DisplayName.

type BroadcastsOn

type BroadcastsOn interface {
	BroadcastOn() []Channel
}

BroadcastsOn is the channels the event goes out on. It is the one method an event must have: Handle refuses an event without it.

type BroadcastsOnConnections

type BroadcastsOnConnections interface {
	BroadcastConnections() []string
}

BroadcastsOnConnections is the connections the event goes out on, which InteractsWithBroadcasting provides.

type BroadcastsVia

type BroadcastsVia interface {
	BroadcastVia(connections ...string) *InteractsWithBroadcasting
}

BroadcastsVia is what PendingBroadcast.Via looks for: an event that can be told which connection to go out on. An event that embeds InteractsWithBroadcasting satisfies it.

type BroadcastsWith

type BroadcastsWith interface {
	BroadcastWith() map[string]any
}

BroadcastsWith is the payload, in place of the event's exported fields.

type Channel

type Channel struct {
	// Name is the channel name, prefix included.
	Name string `json:"name"`
}

Channel is the name an event goes out on.

A private, presence or encrypted channel is this same type with a prefix in front of the name, put there by NewPrivateChannel, NewPresenceChannel or NewEncryptedPrivateChannel. Nothing asks which constructor built a channel; the drivers ask about the prefix.

func NewChannel

func NewChannel(name string) Channel

NewChannel names a channel.

func NewChannelFor

func NewChannelFor(h HasBroadcastChannel) Channel

NewChannelFor names the channel a model is broadcast on.

It is a second constructor rather than one parameter accepting either shape, because Go has no union types.

func NewEncryptedPrivateChannel

func NewEncryptedPrivateChannel(name string) Channel

NewEncryptedPrivateChannel names a private channel whose payloads are encrypted end to end, with the "private-encrypted-" prefix that says so.

The prefix is not decoration: it is how the client knows to decrypt, and how the broadcaster knows to encrypt. A channel that carries the prefix and is not encrypted, or the reverse, fails at the far end with nothing useful said.

func NewPresenceChannel

func NewPresenceChannel(name string) Channel

NewPresenceChannel names a channel that has to be authorized and reports who is listening.

func NewPrivateChannel

func NewPrivateChannel(name string) Channel

NewPrivateChannel names a channel that has to be authorized.

func NewPrivateChannelFor

func NewPrivateChannelFor(h HasBroadcastChannel) Channel

NewPrivateChannelFor names the private channel a model is broadcast on.

func RequestedChannel

func RequestedChannel(name string) (Channel, error)

RequestedChannel is the intake for the one channel name in this package that arrives from outside: the "channel_name" field of the authentication request.

It is a constructor and not a bare string so that the refusal happens once, at the edge, rather than in each driver's Auth. A client that names a tenant is refused rather than trusted or silently re-prefixed: a client that can put "acme:" in front of a channel is a client choosing whose events it hears.

The tenant is added afterwards, by TenantChannel, out of the Grant the Policy issued.

func (Channel) JSONSerialize

func (c Channel) JSONSerialize() any

JSONSerialize answers ToArray, so the two encodings of a channel cannot disagree.

func (Channel) String

func (c Channel) String() string

String is the channel name, and it makes a Channel a fmt.Stringer.

func (Channel) ToArray

func (c Channel) ToArray() map[string]any

ToArray is the channel as it travels inside a broadcast payload: one key, "name".

It is stated here rather than left to whichever driver serializes first.

type Config

type Config struct {
	// Default is the connection used when no name is given, and what
	// [BroadcastManager.GetDefaultDriver] answers.
	Default string
	// Connections is every configured connection, by name.
	Connections map[string]ConnectionConfig
}

Config is the broadcasting configuration a manager is built with.

type ConnectionConfig

type ConnectionConfig struct {
	// Driver is "log", "null", "redis", or a name registered with
	// [BroadcastManager.Extend].
	Driver string
	// Connection is the named connection the redis driver publishes through.
	Connection string
	// Prefix is the key prefix the redis driver puts in front of every channel
	// it publishes on.
	Prefix string
	// Options is what a driver registered with [BroadcastManager.Extend] reads
	// its own settings out of.
	Options map[string]any
}

ConnectionConfig is one configured broadcast connection.

type Dispatcher

type Dispatcher interface {
	// Dispatch fires the event and answers whatever the listeners returned.
	Dispatch(event any, payload ...any) []any
}

Dispatcher is the little of an event dispatcher that a PendingBroadcast uses: it hands the event over, and whatever is listening decides what to do with it.

It is declared here rather than imported from github.com/arandu-io/hesape/events so that an application can broadcast through anything that dispatches -- and so that this package does not pull the outbox and its database in behind it. *events.Dispatcher satisfies it.

type DriverCreator

type DriverCreator func(config ConnectionConfig) (Broadcaster, error)

DriverCreator builds a driver from its configuration. It is what BroadcastManager.Extend registers.

type ExcludesCurrentUser

type ExcludesCurrentUser interface {
	DontBroadcastToCurrentUser(socket string) *InteractsWithSockets
}

ExcludesCurrentUser is what PendingBroadcast.ToOthers looks for: an event that can be told which connection not to go out to. An event that embeds InteractsWithSockets satisfies it.

type Factory

type Factory interface {
	// Connection is the driver registered under a name, or the default one when
	// the name is empty.
	Connection(name string) (Broadcaster, error)
}

Factory is the one method BroadcastEvent.Handle needs from the manager.

It is an interface rather than *BroadcastManager so a test can hand Handle a driver directly.

type FakePendingBroadcast

type FakePendingBroadcast struct {
	PendingBroadcast
}

FakePendingBroadcast is the same three methods, and nothing leaves.

Go has no virtual dispatch through an embedded struct, so the embedding here is for the field set and the substitution is by type: a caller holding a *FakePendingBroadcast calls these methods, and they are the empty ones.

func NewFakePendingBroadcast

func NewFakePendingBroadcast() *FakePendingBroadcast

NewFakePendingBroadcast builds a broadcast that sends nothing.

func (*FakePendingBroadcast) Send

func (f *FakePendingBroadcast) Send() []any

Send dispatches nothing.

func (*FakePendingBroadcast) ToOthers

func (f *FakePendingBroadcast) ToOthers(socket string) *FakePendingBroadcast

ToOthers changes nothing.

func (*FakePendingBroadcast) Via

func (f *FakePendingBroadcast) Via(connection string) *FakePendingBroadcast

Via changes nothing.

type HandlesBroadcastFailure

type HandlesBroadcastFailure interface {
	Failed(ctx context.Context, cause error) error
}

HandlesBroadcastFailure is what the event does when the job carrying it fails.

type HasBackoff

type HasBackoff interface {
	Backoff() time.Duration
}

HasBackoff fills BroadcastEvent.Backoff.

type HasBroadcastChannel

type HasBroadcastChannel interface {
	// BroadcastChannel is the channel name for this instance, e.g. "orders.17".
	BroadcastChannel() string
	// BroadcastChannelRoute is the pattern the channel is registered under,
	// e.g. "orders.{orderId}".
	BroadcastChannelRoute() string
}

HasBroadcastChannel is a model that knows the channel it is broadcast on, and the pattern that channel is authorized under.

NewChannelFor and the other For constructors accept one in place of a name.

type HasBroadcastMiddleware

type HasBroadcastMiddleware interface {
	Middleware() []any
}

HasBroadcastMiddleware is the job middleware of the underlying event.

type HasBroadcastQueue

type HasBroadcastQueue interface {
	BroadcastQueue() string
}

HasBroadcastQueue is the queue the event asks to be pushed onto.

type HasDeleteWhenMissingModels

type HasDeleteWhenMissingModels interface {
	DeleteWhenMissingModels() bool
}

HasDeleteWhenMissingModels fills BroadcastEvent.DeleteWhenMissingModels.

type HasMaxExceptions

type HasMaxExceptions interface {
	MaxExceptions() int
}

HasMaxExceptions fills BroadcastEvent.MaxExceptions.

type HasTimeout

type HasTimeout interface {
	Timeout() time.Duration
}

HasTimeout fills BroadcastEvent.Timeout.

type HasTries

type HasTries interface {
	Tries() int
}

HasTries fills BroadcastEvent.Tries.

type HasUniqueFor

type HasUniqueFor interface {
	UniqueFor() time.Duration
}

HasUniqueFor fills UniqueBroadcastEvent.UniqueFor.

type HasUniqueID

type HasUniqueID interface {
	UniqueID() string
}

HasUniqueID fills UniqueBroadcastEvent.UniqueID.

type InteractsWithBroadcasting

type InteractsWithBroadcasting struct {
	// contains filtered or unexported fields
}

InteractsWithBroadcasting is embedded by an event that chooses the connection it goes out on, and promotes the two methods below onto it.

The zero value is one connection, the default one: [BroadcastConnections] answers a list of one empty name, and an empty name is what BroadcastManager.Driver resolves to the default driver.

func (*InteractsWithBroadcasting) BroadcastConnections

func (i *InteractsWithBroadcasting) BroadcastConnections() []string

BroadcastConnections is the connections the event should be broadcast on.

It never answers an empty list: an event that named no connection is broadcast once, on the default one.

func (*InteractsWithBroadcasting) BroadcastVia

func (i *InteractsWithBroadcasting) BroadcastVia(connections ...string) *InteractsWithBroadcasting

BroadcastVia names the connections the event goes out on. Passing none restores the default connection.

type InteractsWithSockets

type InteractsWithSockets struct {
	// Socket is the socket id of the connection that raised the event. Empty
	// means the event goes to everyone.
	//
	// It is tagged so that it lands in a broadcast payload under "socket"
	// rather than "Socket".
	Socket string `json:"socket"`
}

InteractsWithSockets is embedded by an event that can exclude the connection that raised it, and promotes the two methods below onto it:

type OrderShipped struct {
	broadcasting.InteractsWithSockets
	OrderID string
}

The socket id is passed in rather than read off an ambient request, because a request in Go is a value a handler holds. BroadcastManager.Socket is where it comes from, off the X-Socket-ID header.

func (*InteractsWithSockets) BroadcastToEveryone

func (i *InteractsWithSockets) BroadcastToEveryone() *InteractsWithSockets

BroadcastToEveryone clears the exclusion, so the event reaches every connection.

func (*InteractsWithSockets) DontBroadcastToCurrentUser

func (i *InteractsWithSockets) DontBroadcastToCurrentUser(socket string) *InteractsWithSockets

DontBroadcastToCurrentUser excludes the connection that raised the event from receiving it.

type PendingBroadcast

type PendingBroadcast struct {
	// contains filtered or unexported fields
}

PendingBroadcast is an event on its way to the dispatcher, with two things still adjustable.

Nothing sends it on the way out of scope: PendingBroadcast.Send is called by hand, at the end of the chain.

manager.Event(e).ToOthers(socket).Send()

func NewPendingBroadcast

func NewPendingBroadcast(events Dispatcher, event any) *PendingBroadcast

NewPendingBroadcast builds the broadcast an event is sent through.

func (*PendingBroadcast) Event

func (p *PendingBroadcast) Event() any

Event exposes the event this broadcast is carrying, which is the only way to see what is about to be sent.

func (*PendingBroadcast) Send

func (p *PendingBroadcast) Send() []any

Send dispatches the event, and answers what the listeners answered.

Calling it twice dispatches twice: it is an ordinary method, and nothing marks a broadcast as already sent.

func (*PendingBroadcast) ToOthers

func (p *PendingBroadcast) ToOthers(socket string) *PendingBroadcast

ToOthers broadcasts the event to everyone except the current user.

The socket id is an argument because there is no ambient request to read it off -- see InteractsWithSockets. It comes from BroadcastManager.Socket.

func (*PendingBroadcast) Via

func (p *PendingBroadcast) Via(connection string) *PendingBroadcast

Via broadcasts the event using a specific broadcaster.

An event that cannot be told which connection to use -- one that does not embed InteractsWithBroadcasting -- is left alone rather than refused. An empty connection is the default one.

type Queue

type Queue interface {
	// PushOn puts the job on a named queue of a named connection. Both names
	// may be empty: the default queue of the default connection.
	//
	// The job is *BroadcastEvent, or *UniqueBroadcastEvent when the event asked
	// to be unique.
	PushOn(ctx context.Context, g auth.Grant, connection, queue string, job any) error
}

Queue is the little of a queue factory that BroadcastManager.Queue uses.

It is declared here rather than imported from github.com/arandu-io/hesape/queue so that an application can broadcast without a queue behind it, and so that this package does not pull the worker, the drivers and their database in.

type Router

type Router interface {
	// Handle registers a handler, e.g. Handle("GET /broadcasting/auth", h).
	Handle(pattern string, handler http.Handler)
}

Router is the little of a router that BroadcastManager.Routes and BroadcastManager.UserRoutes use: a pattern and a handler.

It is declared here rather than imported from github.com/arandu-io/hesape/routing so that this package does not depend on the router to be usable, and because *net/http.ServeMux satisfies it as it stands -- its patterns carry the method as well as the path.

type ShouldBroadcastNow

type ShouldBroadcastNow interface {
	ShouldBroadcastNow() bool
}

ShouldBroadcastNow returning true makes the event skip the queue: it is published in the goroutine that raised it.

type ShouldRescue

type ShouldRescue interface {
	ShouldRescue() bool
}

ShouldRescue returning true swallows a failure to publish rather than returning it.

type UniqueBroadcastEvent

type UniqueBroadcastEvent struct {
	BroadcastEvent

	// UniqueID is the lock identifier.
	UniqueID string
	// UniqueFor is how long the lock is held.
	UniqueFor time.Duration
}

UniqueBroadcastEvent is a BroadcastEvent that will not be queued twice while one is still in flight.

BroadcastManager.Queue is what reads the two fields: it takes a lock under UniqueID before pushing, and drops the event when somebody already holds it.

func NewUniqueBroadcastEvent

func NewUniqueBroadcastEvent(event any) *UniqueBroadcastEvent

NewUniqueBroadcastEvent builds the job for an event that asked to be unique.

type UniqueLock

type UniqueLock interface {
	// Acquire is true when this caller took the lock, false when somebody
	// already holds it.
	Acquire(ctx context.Context, g auth.Grant, key string, expiresAfter time.Duration) (bool, error)
}

UniqueLock is the little of a lock that BroadcastManager.Queue uses to keep a unique event from being pushed twice.

It is an interface rather than github.com/arandu-io/hesape/bus.UniqueLock so that a broadcast does not require a cache to be configured, and so this package does not depend on the bus. bus.UniqueLock reaches it through the two-line adapter that turns the key into its UniqueJob.

type UserResolver

type UserResolver interface {
	// ResolveAuthenticatedUser is the user payload the socket server is given
	// for the connection, or nil when no callback was registered.
	ResolveAuthenticatedUser(ctx context.Context, r *http.Request) (any, error)
}

UserResolver is the optional fourth method of a driver: it answers who the connection belongs to.

BroadcastController.AuthenticateUser asks the driver for it, and refuses with 403 when the driver does not have it.

Directories

Path Synopsis
Package broadcasters holds the drivers a broadcast is published through.
Package broadcasters holds the drivers a broadcast is published through.

Jump to

Keyboard shortcuts

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