Documentation
¶
Overview ¶
Package realtime implements the Nimbu realtime protocol: grant-authenticated ActionCable sessions, live-query validation and a reconnecting watch loop.
Index ¶
- Constants
- Variables
- func ParseMessage(raw json.RawMessage) (*Control, *Event, error)
- func ValidateQuery(q map[string]string) error
- type BackoffOptions
- type Conn
- type Control
- type Deduper
- type DialOptions
- type Event
- type ExitReason
- type Frame
- type GrantRejected
- type Handler
- type Identifier
- type Lease
- type QueryError
- type SessionOptions
- type SessionResult
- type Subscribed
- type SubscriptionError
- type WatchOptions
Constants ¶
const ( FrameWelcome = "welcome" FramePing = "ping" FrameConfirmSubscription = "confirm_subscription" FrameRejectSubscription = "reject_subscription" FrameDisconnect = "disconnect" )
Frame types sent by the ActionCable server.
const ( ControlSubscribed = "subscribed" ControlRenewed = "renewed" ControlUnsubscribed = "unsubscribed" ControlSubscriptionErr = "subscription_error" ControlConfirmSubscribe = FrameConfirmSubscription )
Control message types carried inside a frame's message payload.
const ( EventAdded = "added" EventChanged = "changed" EventRemoved = "removed" EventResync = "resync" )
Event names carried by realtime events.
const ( DefaultMaxAttempts = 10 DefaultBackoffBase = time.Second DefaultBackoffMax = 30 * time.Second )
Defaults for the reconnect loop.
const ChannelName = "RealtimeApiChannel"
ChannelName is the ActionCable channel serving live queries.
const DefaultDedupeCapacity = 4096
DefaultDedupeCapacity is the number of event ids remembered per watch.
const DefaultPingTimeout = 30 * time.Second
DefaultPingTimeout is how long a session waits for any frame before it treats the socket as dead.
const ResourceChannelEntries = "channel_entries"
ResourceChannelEntries is the realtime resource for channel entries.
const Subprotocol = "actioncable-v1-json"
Subprotocol is the ActionCable JSON subprotocol used by the realtime socket.
Variables ¶
var ErrPingTimeout = errors.New("realtime: no frames received before the ping timeout")
ErrPingTimeout is returned when no frame arrives within the ping timeout.
Functions ¶
func ParseMessage ¶
func ParseMessage(raw json.RawMessage) (*Control, *Event, error)
ParseMessage decodes a frame message into either a control or an event.
func ValidateQuery ¶
ValidateQuery checks a live-query filter map against the rules the server applies, so the CLI can fail with a precise message before connecting.
Types ¶
type BackoffOptions ¶
type BackoffOptions struct {
Base time.Duration
Max time.Duration
Jitter func(time.Duration) time.Duration
}
BackoffOptions tunes the delay between reconnect attempts.
type Conn ¶
type Conn interface {
Read(ctx context.Context) ([]byte, error)
Write(ctx context.Context, b []byte) error
Close(reason string) error
}
Conn is the minimal websocket surface the session loop needs.
type Control ¶
type Control struct {
Type string
Code string
Subscribed *Subscribed
Inactive bool
Raw json.RawMessage
}
Control is a non-event message from the realtime channel.
type Deduper ¶
type Deduper struct {
// contains filtered or unexported fields
}
Deduper remembers recently seen event ids in a bounded FIFO set. Delivery is at-least-once, so repeated ids must be dropped.
func NewDeduper ¶
NewDeduper creates a Deduper holding at most capacity ids.
type DialOptions ¶
DialOptions tunes the websocket handshake.
type Event ¶
type Event struct {
EventID string `json:"event_id"`
Event string `json:"event"`
Resource string `json:"resource"`
ParentID string `json:"parent_id"`
ID string `json:"id"`
Type string `json:"type"`
Object json.RawMessage `json:"object,omitempty"`
Changeset json.RawMessage `json:"changeset,omitempty"`
OccurredAt string `json:"occurred_at"`
Revision int64 `json:"revision"`
Raw json.RawMessage `json:"-"`
}
Event is a single live-query event envelope.
type ExitReason ¶
type ExitReason int
ExitReason explains why Watch returned.
const ( // ExitInterrupted means the caller's context was cancelled. ExitInterrupted ExitReason = iota // ExitTimeout means the configured watch duration elapsed. ExitTimeout // ExitOnce means the first event arrived and Once was set. ExitOnce // ExitFatal means watching stopped on an error that will not resolve. ExitFatal )
func Watch ¶
func Watch(ctx context.Context, opts WatchOptions) (ExitReason, error)
Watch mints grants, dials and runs sessions until the context ends, the first event arrives under Once, or the attempt budget runs out.
func (ExitReason) String ¶
func (r ExitReason) String() string
String renders the exit reason for logs.
type Frame ¶
type Frame struct {
Type string `json:"type"`
Identifier string `json:"identifier"`
Message json.RawMessage `json:"message"`
Reason string `json:"reason"`
Reconnect *bool `json:"reconnect"`
}
Frame is one decoded ActionCable frame.
func ParseFrame ¶
ParseFrame decodes one ActionCable frame.
type GrantRejected ¶
type GrantRejected struct {
Reason string
}
GrantRejected is returned when the server refuses the grant, which means a fresh grant must be minted before reconnecting.
func (*GrantRejected) Error ¶
func (e *GrantRejected) Error() string
type Identifier ¶
type Identifier struct {
Channel string `json:"channel"`
Resource string `json:"resource"`
ParentID string `json:"parent_id"`
Query map[string]string `json:"query,omitempty"`
}
Identifier is the ActionCable subscription identifier payload.
func NewIdentifier ¶
func NewIdentifier(parentID string, query map[string]string) Identifier
NewIdentifier builds a channel-entries identifier for a parent channel.
func (Identifier) String ¶
func (i Identifier) String() (string, error)
String returns the canonical JSON encoding used as the ActionCable identifier.
type QueryError ¶
QueryError reports a live-query filter that the server would reject.
func (*QueryError) Error ¶
func (e *QueryError) Error() string
type SessionOptions ¶
type SessionOptions struct {
Identifier Identifier
PingTimeout time.Duration
Now func() time.Time
Dedupe *Deduper
Handler Handler
}
SessionOptions configures one socket lifetime.
type SessionResult ¶
type SessionResult int
SessionResult explains why a single socket lifetime ended.
const ( // ResultInactive means the subscription lapsed and must be recreated. ResultInactive SessionResult = iota // ResultDisconnect means the socket dropped and may be retried. ResultDisconnect // ResultOnce means the caller asked to stop after the first event. ResultOnce // ResultDone means the context was cancelled. ResultDone )
func RunSession ¶
func RunSession(ctx context.Context, conn Conn, opts SessionOptions) (SessionResult, error)
RunSession drives one socket from welcome to teardown.
func (SessionResult) String ¶
func (r SessionResult) String() string
String renders the result for logs.
type Subscribed ¶
type Subscribed struct {
SubscriptionID string
Resource string
Firehose bool
Query map[string]string
SeqAtCreate int64
Lease Lease
Resync bool
Raw json.RawMessage
}
Subscribed describes an accepted subscription and its lease.
type SubscriptionError ¶
type SubscriptionError struct {
Code string
}
SubscriptionError is a subscription_error control returned by the server.
func (*SubscriptionError) Error ¶
func (e *SubscriptionError) Error() string
func (*SubscriptionError) Fatal ¶
func (e *SubscriptionError) Fatal() bool
Fatal reports whether retrying the same subscription is pointless.
type WatchOptions ¶
type WatchOptions struct {
Identifier Identifier
MintGrant func(ctx context.Context) (string, error)
Dial func(ctx context.Context, grant string) (Conn, error)
Handler Handler
Once bool
Timeout time.Duration
PingTimeout time.Duration
MaxAttempts int
Backoff BackoffOptions
Now func() time.Time
Sleep func(ctx context.Context, d time.Duration) error
}
WatchOptions configures the reconnecting watch loop.