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
- Variables
- func CutTenant(name string) (tenant, channel string, found bool)
- func TenantChannel(g auth.Grant, c Channel) (string, error)
- func TenantChannels(g auth.Grant, channels []Channel) ([]string, error)
- type AnonymousEvent
- func (e *AnonymousEvent) As(name string) *AnonymousEvent
- func (e *AnonymousEvent) BroadcastAs() string
- func (e *AnonymousEvent) BroadcastOn() []Channel
- func (e *AnonymousEvent) BroadcastWith() map[string]any
- func (e *AnonymousEvent) Send() []any
- func (e *AnonymousEvent) SendNow() []any
- func (e *AnonymousEvent) ShouldBroadcastNow() bool
- func (e *AnonymousEvent) ToOthers(socket string) *AnonymousEvent
- func (e *AnonymousEvent) Via(connection string) *AnonymousEvent
- func (e *AnonymousEvent) With(payload any) *AnonymousEvent
- type Arrayable
- type BroadcastController
- type BroadcastError
- type BroadcastEvent
- func (b *BroadcastEvent) Clone() *BroadcastEvent
- func (b *BroadcastEvent) DisplayName() string
- func (b *BroadcastEvent) Failed(ctx context.Context, cause error) error
- func (b *BroadcastEvent) Handle(ctx context.Context, g auth.Grant, manager Factory) error
- func (b *BroadcastEvent) Middleware() []any
- type BroadcastManager
- func (m *BroadcastManager) ChannelRoutes(r Router)
- func (m *BroadcastManager) Connection(name string) (Broadcaster, error)
- func (m *BroadcastManager) Driver(name string) (Broadcaster, error)
- func (m *BroadcastManager) Event(event any) *PendingBroadcast
- func (m *BroadcastManager) Extend(driver string, creator DriverCreator) *BroadcastManager
- func (m *BroadcastManager) ForgetDrivers() *BroadcastManager
- func (m *BroadcastManager) GetDefaultDriver() string
- func (m *BroadcastManager) On(channels ...Channel) *AnonymousEvent
- func (m *BroadcastManager) Presence(channel string) *AnonymousEvent
- func (m *BroadcastManager) Private(channel string) *AnonymousEvent
- func (m *BroadcastManager) Purge(name string)
- func (m *BroadcastManager) Queue(ctx context.Context, g auth.Grant, event any) error
- func (m *BroadcastManager) Routes(r Router)
- func (m *BroadcastManager) SetDefaultDriver(name string)
- func (m *BroadcastManager) Socket(r *http.Request) string
- func (m *BroadcastManager) UserRoutes(r Router)
- type Broadcaster
- type BroadcastsAs
- type BroadcastsOn
- type BroadcastsOnConnections
- type BroadcastsVia
- type BroadcastsWith
- type Channel
- func NewChannel(name string) Channel
- func NewChannelFor(h HasBroadcastChannel) Channel
- func NewEncryptedPrivateChannel(name string) Channel
- func NewPresenceChannel(name string) Channel
- func NewPrivateChannel(name string) Channel
- func NewPrivateChannelFor(h HasBroadcastChannel) Channel
- func RequestedChannel(name string) (Channel, error)
- type Config
- type ConnectionConfig
- type Dispatcher
- type DriverCreator
- type ExcludesCurrentUser
- type Factory
- type FakePendingBroadcast
- type HandlesBroadcastFailure
- type HasBackoff
- type HasBroadcastChannel
- type HasBroadcastMiddleware
- type HasBroadcastQueue
- type HasDeleteWhenMissingModels
- type HasMaxExceptions
- type HasTimeout
- type HasTries
- type HasUniqueFor
- type HasUniqueID
- type InteractsWithBroadcasting
- type InteractsWithSockets
- type PendingBroadcast
- type Queue
- type Router
- type ShouldBroadcastNow
- type ShouldRescue
- type UniqueBroadcastEvent
- type UniqueLock
- type UserResolver
Constants ¶
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.
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.
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.
const ChannelNameField = "channel_name"
ChannelNameField is the field the socket client sends the channel it wants to listen on under.
const SocketIDHeader = "X-Socket-ID"
SocketIDHeader is the header BroadcastManager.Socket reads the socket id of the calling connection off.
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 ¶
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.
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 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 ¶
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 ¶
NewPresenceChannel names a channel that has to be authorized and reports who is listening.
func NewPrivateChannel ¶
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 ¶
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 ¶
JSONSerialize answers ToArray, so the two encodings of a channel cannot disagree.
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 ¶
HandlesBroadcastFailure is what the event does when the job carrying it fails.
type HasBackoff ¶
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 ¶
HasTimeout fills BroadcastEvent.Timeout.
type HasUniqueFor ¶
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.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package broadcasters holds the drivers a broadcast is published through.
|
Package broadcasters holds the drivers a broadcast is published through. |