uplink

package
v0.14.0 Latest Latest
Warning

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

Go to latest
Published: Aug 27, 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.

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 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) *Client

NewClient creates a new WebSocket client.

The returned client already handles the two exchanges that belong to the link rather than to any tenant: a clock sync and an organization lookup, both issued on every connect. Tenants layer their own handlers 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) OnConnect

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

OnConnect registers a callback invoked each time a WebSocket connection is established. The handler can use Send to queue messages for the new connection.

Callbacks accumulate rather than replace: the connector registers its own for time sync and org info, and feature reporters add theirs on top. They run in registration order on the connection goroutine, so a callback that needs to do real work should hand off to its own.

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) 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 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 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