realtime

package
v0.4.1 Latest Latest
Warning

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

Go to latest
Published: Sep 30, 2026 License: MIT Imports: 14 Imported by: 0

Documentation

Overview

Package realtime implements the Nimbu realtime protocol: grant-authenticated ActionCable sessions, live-query validation and a reconnecting watch loop.

Index

Constants

View Source
const (
	FrameWelcome             = "welcome"
	FramePing                = "ping"
	FrameConfirmSubscription = "confirm_subscription"
	FrameRejectSubscription  = "reject_subscription"
	FrameDisconnect          = "disconnect"
)

Frame types sent by the ActionCable server.

View Source
const (
	ControlSubscribed       = "subscribed"
	ControlRenewed          = "renewed"
	ControlUnsubscribed     = "unsubscribed"
	ControlSubscriptionErr  = "subscription_error"
	ControlConfirmSubscribe = FrameConfirmSubscription
)

Control message types carried inside a frame's message payload.

View Source
const (
	EventAdded   = "added"
	EventChanged = "changed"
	EventRemoved = "removed"
	EventResync  = "resync"
)

Event names carried by realtime events.

View Source
const (
	DefaultMaxAttempts = 10
	DefaultBackoffBase = time.Second
	DefaultBackoffMax  = 30 * time.Second
)

Defaults for the reconnect loop.

View Source
const ChannelName = "RealtimeApiChannel"

ChannelName is the ActionCable channel serving live queries.

View Source
const DefaultDedupeCapacity = 4096

DefaultDedupeCapacity is the number of event ids remembered per watch.

View Source
const DefaultPingTimeout = 30 * time.Second

DefaultPingTimeout is how long a session waits for any frame before it treats the socket as dead.

View Source
const ResourceChannelEntries = "channel_entries"

ResourceChannelEntries is the realtime resource for channel entries.

View Source
const Subprotocol = "actioncable-v1-json"

Subprotocol is the ActionCable JSON subprotocol used by the realtime socket.

Variables

View Source
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

func ValidateQuery(q map[string]string) error

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.

func Dial

func Dial(ctx context.Context, siteHost, grant string, opts DialOptions) (Conn, error)

Dial opens the realtime websocket for a site using a one-time grant. The grant is carried in the query string and is scrubbed from every error.

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

func NewDeduper(capacity int) *Deduper

NewDeduper creates a Deduper holding at most capacity ids.

func (*Deduper) Seen

func (d *Deduper) Seen(id string) bool

Seen records an id and reports whether it was already known. An empty id is never treated as a duplicate.

type DialOptions

type DialOptions struct {
	Insecure   bool
	HTTPClient *http.Client
	Headers    http.Header
}

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

func ParseFrame(b []byte) (Frame, error)

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 Handler

type Handler interface {
	OnEvent(ev Event)
	OnControl(c Control)
	OnStatus(msg string)
}

Handler receives everything a session observes.

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 Lease

type Lease struct {
	ExpiresAt  float64
	RenewAfter float64
}

Lease describes when a subscription must be renewed, in unix seconds.

type QueryError

type QueryError struct {
	Key    string
	Reason string
}

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.

Jump to

Keyboard shortcuts

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