uplink

package
v0.16.2 Latest Latest
Warning

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

Go to latest
Published: Sep 29, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

Documentation

Overview

Package uplink implements the cluster's control-plane link to Miren Cloud: exactly one persistent, authenticated WebSocket per cluster, carrying typed coordination messages multiplexed by type.

RFD-94 calls this the unified uplink and treats it as the backbone for all runtime↔cloud coordination. Features are tenants of it — Miren Anywhere for NAT traversal, app state reporting, and whatever follows — and register their own handlers and connect hooks rather than the link knowing anything about them.

The pipe is payload-agnostic on purpose, and should stay that way. It moves opaque Envelopes and understands no apps, deploys, or POPs. That is separate from it being delivery-unreliable, which it also is: the outbox is bounded and drops on overflow, and a reconnect discards whatever was queued. Tenants that need completeness reconcile from their own durable source rather than expecting the wire to have carried everything.

The exception to "knows nothing" is deliberate and small: clock offset and organization identity are established here, because both are properties of the link itself rather than of any one tenant.

Index

Constants

View Source
const (
	TypeTimeRequest     = "time.request"
	TypeTimeResponse    = "time.response"
	TypeOrgInfoRequest  = "org.info.request"
	TypeOrgInfoResponse = "org.info.response"
)

Control-plane message types. Feature message types are namespaced by their own family (app.*, deploy.*) and defined by the tenant that owns them, not here.

View Source
const (
	TypeSessionHello   = "session.hello"
	TypeSessionWelcome = "session.welcome"
	TypeSessionReject  = "session.reject"

	HandshakeVersion1 uint = 1

	CapabilityPopConnect = "pop-connect"
	CapabilityRPCRelay   = "rpc-relay"
	CapabilityEntitySync = "entity-sync"
	CapabilityAppHealth  = "app-health"
	// CapabilityServerLifecycle lets cloud restart and upgrade the server and
	// follow the operations it started.
	CapabilityServerLifecycle = "server-lifecycle"
	// CapabilityClusterNetwork carries how the cluster can be reached: the
	// network half of the legacy status report.
	CapabilityClusterNetwork = "cluster-network"
	// CapabilityClusterResources carries whole-host CPU, memory, and storage
	// utilization: the resource half of the legacy status report.
	CapabilityClusterResources = "cluster-resources"
)

Variables

This section is empty.

Functions

func SpreadOnConnect

func SpreadOnConnect(window time.Duration) time.Duration

SpreadOnConnect returns a delay a tenant should wait before starting the work it does on each connect, so that work is spread across the fleet rather than landing as one spike.

Exported rather than kept inside Run because jittering the reconnect alone only moves the spike: a fleet that comes back staggered but then has every cluster immediately stream a snapshot has changed when the pile-up happens, not whether it does. Tenants need the same treatment for their own connect work.

The window is the caller's choice because it depends on what the work is worth delaying. Anything on the ephemeral tier can afford a generous one: nothing downstream distinguishes a snapshot that starts now from one that starts a minute from now.

Types

type CapabilityOffer added in v0.15.0

type CapabilityOffer struct {
	Name     string          `json:"name"`
	Versions []uint          `json:"versions"`
	Offer    json.RawMessage `json:"offer,omitempty"`
}

CapabilityOffer describes one protocol family the runtime can speak. The session layer negotiates the version but leaves Offer to the capability. For example, entity sync will use it to offer export-schema digests without teaching the uplink what an entity is.

type CapabilityOfferFunc added in v0.15.0

type CapabilityOfferFunc func(context.Context) (json.RawMessage, bool)

CapabilityOfferFunc builds the connection-specific payload for an offered capability. Returning false omits the capability from that connection.

type CapabilitySelection added in v0.15.0

type CapabilitySelection struct {
	Name    string          `json:"name"`
	Version uint            `json:"version"`
	Config  json.RawMessage `json:"config,omitempty"`
}

CapabilitySelection is cloud's choice for one offered capability. Config is owned and decoded by that capability, not by the session layer.

type Client

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

Client maintains a persistent WebSocket connection to the cloud coordination service with automatic reconnection.

func NewClient

func NewClient(cloudURL string, authClient *cloudauth.AuthClient, router *MessageRouter, log *slog.Logger, opts ...ClientOption) *Client

NewClient creates a new WebSocket client.

Without WithSession, the client preserves the legacy clock-sync and organization-lookup requests on every connection. With it, those link facts arrive in session.welcome instead. Tenants layer their own handlers, capability offers, and hooks on top.

func (*Client) Handle

func (c *Client) Handle(msgType string, handler MessageHandler)

Handle registers a handler for an inbound message type, which is how a feature becomes a tenant of the link. Message types are namespaced per family (app.*, deploy.*) so the link stays a shared pipe rather than accumulating per-feature special cases.

func (*Client) OfferCapability added in v0.15.0

func (c *Client) OfferCapability(offer CapabilityOffer)

OfferCapability adds a protocol family to the next session hello. Offers are snapshotted for each connection, so tenants should register before Run.

func (*Client) OfferCapabilityFunc added in v0.15.0

func (c *Client) OfferCapabilityFunc(name string, versions []uint, provide CapabilityOfferFunc)

OfferCapabilityFunc adds a capability whose offer payload is rebuilt before each connection's session handshake. Returning false omits only that capability from the connection, so temporary source failures do not prevent other tenants from negotiating.

func (*Client) OnConnect

func (c *Client) OnConnect(fn func(ctx context.Context))

OnConnect registers a callback invoked each time the link is ready. For a negotiated connection that means after session.welcome has been validated; for a legacy connection it means immediately after the WebSocket connects. The handler can use Send to queue messages for the new connection.

Callbacks accumulate rather than replace, so each tenant can add its own. They run in registration order on the connection goroutine, so a callback that needs to do real work should hand off to its own goroutine.

The context is scoped to the connection, so work handed off that way is cancelled when the connection drops. That is the right lifetime for it: anything still queued at that point is discarded on reconnect anyway, so finishing would only produce a message nobody sends.

func (*Client) OnSession added in v0.15.0

func (c *Client) OnSession(fn func(ctx context.Context, session Session))

OnSession registers a callback invoked after cloud's welcome has been validated. It is never called for a legacy connection. The context ends with the negotiated session, and the Session value does not change underneath the callback. Callbacks run in registration order before the read and write loops start, so any real work should be handed off to a goroutine.

func (*Client) Run

func (c *Client) Run(ctx context.Context) error

Run maintains the WebSocket connection with reconnection. It blocks until the context is cancelled.

func (*Client) Send

func (c *Client) Send(env *Envelope)

Send queues an envelope for delivery to the cloud. Non-blocking; drops the message if the outbox is full.

func (*Client) SendBlocking

func (c *Client) SendBlocking(ctx context.Context, env *Envelope) error

SendBlocking queues an envelope, waiting for room instead of dropping when the outbox is full. It returns when the envelope is queued, or with the context's error if the connection goes away first.

Send's drop-on-overflow is the right behavior for a single small message whose loss self-heals, but wrong for a stream. A tenant sending a bounded sequence — a snapshot in batches, a backfill walking a watermark — is pushing faster than one message per connect, and silent drops there don't read as "we lost a sample," they read as "that batch never existed." For a snapshot that is indistinguishable from apps having been deleted.

Blocking makes the write loop's drain rate the natural throttle, which is also what lets two streaming tenants share the link without coordinating: neither can starve the other by filling the outbox, because filling it just slows the filler down. Callers must pass the connection-scoped context so a sender parked here is released when the connection drops rather than waking up to write into a socket that is gone.

func (*Client) SendContext

func (c *Client) SendContext(ctx context.Context, env *Envelope) error

SendContext queues an envelope and waits for room rather than dropping, for a tenant that cannot treat loss as a self-healing condition. The best-effort Send above remains the right call for state a later report repairs; this one is for a stream whose gap nothing downstream can reconstruct.

What it promises is narrow: the envelope reached the outbox of the connection live at the time, so it will be written unless that connection dies first. It is not an acknowledgement, and a reconnect discards whatever is still queued (see drainOutbox). A tenant that needs more than "delivered or the link broke" has to notice the break — which, for anything session-shaped, means tearing the session down and letting the caller retry.

func (*Client) SendMessage

func (c *Client) SendMessage(msgType string, data any) error

SendMessage marshals data and queues it for delivery.

func (*Client) SendMessageBlocking

func (c *Client) SendMessageBlocking(ctx context.Context, msgType string, data any) error

SendMessageBlocking marshals data and queues it, waiting for outbox room. See SendBlocking for when to prefer this over SendMessage.

type ClientOption added in v0.15.0

type ClientOption func(*Client)

ClientOption changes how a Client establishes its connections.

func WithSession added in v0.15.0

func WithSession(identity SessionIdentity) ClientOption

WithSession enables the negotiated session handshake. Without it the client speaks the legacy bootstrap, which cloud still accepts from older runtimes.

func WithStatus added in v0.15.0

func WithStatus(report StatusFunc) ClientOption

WithStatus reports connection lifecycle changes for operator diagnostics.

type Envelope

type Envelope struct {
	Type string          `json:"type"`
	Data json.RawMessage `json:"data"`
}

Envelope is the wire format for all messages on the WebSocket.

type MessageHandler

type MessageHandler func(ctx context.Context, data json.RawMessage) error

MessageHandler processes an inbound message.

type MessageRouter

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

MessageRouter dispatches inbound messages by type.

func NewMessageRouter

func NewMessageRouter() *MessageRouter

NewMessageRouter creates a new MessageRouter.

func (*MessageRouter) Dispatch

func (r *MessageRouter) Dispatch(ctx context.Context, env Envelope) error

Dispatch routes a message to the appropriate handler.

func (*MessageRouter) Handle

func (r *MessageRouter) Handle(msgType string, handler MessageHandler)

Handle registers a handler for a given message type.

type OrgInfoResponse

type OrgInfoResponse struct {
	OrganizationID string `json:"organization_id"`
}

OrgInfoResponse is the payload for org.info.response messages.

type Session added in v0.15.0

type Session struct {
	ID               string
	HandshakeVersion uint
	RuntimeVersion   string
	OrganizationID   string
	ClockOffset      time.Duration
	Capabilities     []CapabilitySelection
	// IdentityIssuerURL is the workload identity anchor cloud named in the
	// welcome, or empty. See SessionWelcome.
	IdentityIssuerURL string
}

Session is immutable negotiated state scoped to one WebSocket connection.

func (Session) Capability added in v0.15.0

func (s Session) Capability(name string) (CapabilitySelection, bool)

Capability returns the selected capability with the given name.

type SessionHello added in v0.15.0

type SessionHello struct {
	HandshakeVersions []uint `json:"handshake_versions"`
	RuntimeVersion    string `json:"runtime_version"`
	// RuntimeInstanceID lets cloud tell a reconnect from a restart. Optional
	// so older runtimes stay decodable.
	RuntimeInstanceID string `json:"runtime_instance_id,omitempty"`
	// RuntimeCommit and RuntimeBuildDate identify the build exactly.
	// RuntimeVersion carries only a short sha, which is ambiguous once git
	// lengthens abbreviations under a collision and carries no order, so
	// cloud needs the full commit to match a build against its channel and
	// the build date to say which of two main builds is newer. Both are
	// absent from a binary built without hack/build.sh's ldflags.
	RuntimeCommit    string            `json:"runtime_commit,omitempty"`
	RuntimeBuildDate time.Time         `json:"runtime_build_date,omitzero"`
	ClientTime       time.Time         `json:"client_time"`
	Capabilities     []CapabilityOffer `json:"capabilities"`
}

SessionHello is the first envelope sent on a negotiated uplink connection. The bootstrap shape is deliberately small: future protocol families evolve behind capabilities rather than adding their state directly here.

type SessionIdentity added in v0.16.0

type SessionIdentity struct {
	RuntimeVersion    string
	RuntimeInstanceID string
	// RuntimeCommit and RuntimeBuildDate are optional; see SessionHello.
	RuntimeCommit    string
	RuntimeBuildDate time.Time
}

type SessionReject added in v0.15.0

type SessionReject struct {
	Reason                     string `json:"reason"`
	SupportedHandshakeVersions []uint `json:"supported_handshake_versions,omitempty"`
}

SessionReject explains why cloud could not establish a negotiated session. It is a bootstrap message, so it must remain decodable by every handshake version.

type SessionWelcome added in v0.15.0

type SessionWelcome struct {
	HandshakeVersion   uint                  `json:"handshake_version"`
	SessionID          string                `json:"session_id"`
	OrganizationID     string                `json:"organization_id"`
	ServerReceiveTime  time.Time             `json:"server_receive_time"`
	ServerTransmitTime time.Time             `json:"server_transmit_time"`
	Capabilities       []CapabilitySelection `json:"capabilities"`
	// IdentityIssuerURL is where cloud anchors this cluster's workload
	// identity. The status poll's response carried it on every report so a
	// cluster registered before anchors existed could learn its own; the
	// welcome carries it on every session for the same reason. Empty when
	// cloud is not serving discovery, or from a cloud that predates it.
	IdentityIssuerURL string `json:"identity_issuer_url,omitempty"`
}

SessionWelcome establishes the session and selects the protocol families both peers will use for this connection.

type Status added in v0.15.0

type Status struct {
	State     string
	Session   *Session
	LastError string
	RetryAt   time.Time
}

Status describes the current state of the reconnecting uplink client.

type StatusFunc added in v0.15.0

type StatusFunc func(Status)

StatusFunc receives uplink lifecycle changes. Callbacks must return quickly.

type TimeRequest

type TimeRequest struct {
	ClientTransmitTime time.Time `json:"client_transmit_time"`
}

TimeRequest is the payload for time.request messages.

type TimeResponse

type TimeResponse struct {
	ClientTransmitTime time.Time `json:"client_transmit_time"`
	ServerReceiveTime  time.Time `json:"server_receive_time"`
	ServerTransmitTime time.Time `json:"server_transmit_time"`
}

TimeResponse carries the four timestamps needed for simplified NTP.

Jump to

Keyboard shortcuts

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