cluster

package
v1.17.0 Latest Latest
Warning

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

Go to latest
Published: Aug 3, 2026 License: MIT Imports: 12 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type DefaultResolver added in v1.14.0

type DefaultResolver struct{}

DefaultResolver wraps the system DNS resolver.

func (DefaultResolver) LookupHost added in v1.14.0

func (DefaultResolver) LookupHost(ctx context.Context, host string) ([]string, error)

type Discovery added in v1.14.0

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

Discovery periodically resolves a Kubernetes Headless Service hostname via DNS and reconciles the PeerManager peer list using RCU (atomic swap).

Algorithm:

  1. Every N seconds, resolve host → []IP via net.LookupHost.
  2. Sort + deduplicate the result set.
  3. Diff against the currently known set.
  4. For IPs not in current set → PeerManager.AddPeer.
  5. For IPs in current set but not in resolved set → PeerManager.RemovePeer.

Why DNS over K8s API (client-go):

  • client-go adds ~40MB to the binary (REST client, informers, protobuf).
  • DNS resolution uses only net.LookupHost — zero external dependencies.
  • Headless Services in Kubernetes automatically create SRV/A records for all ready pods, providing the same discovery as API label selectors.
  • Preserves the "single static binary" philosophy.

func NewDiscovery added in v1.14.0

func NewDiscovery(resolver Resolver, cfg DiscoveryConfig, pm *PeerManager, logger *slog.Logger) *Discovery

NewDiscovery creates a Discovery that will periodically resolve hostname and reconcile peers on the given PeerManager.

func (*Discovery) KnownPeers added in v1.14.0

func (d *Discovery) KnownPeers() []string

KnownPeers returns a sorted copy of the currently discovered peer addresses.

func (*Discovery) ResolveOnce added in v1.14.0

func (d *Discovery) ResolveOnce(ctx context.Context)

ResolveOnce performs a single resolution + reconciliation (exported for testing).

func (*Discovery) Start added in v1.14.0

func (d *Discovery) Start(ctx context.Context)

Start begins the background DNS polling loop. It resolves immediately on start, then repeats every d.interval. Cancel ctx to stop.

type DiscoveryConfig added in v1.14.0

type DiscoveryConfig struct {
	Enabled  bool          // whether discovery is active
	Host     string        // headless service FQDN e.g. "aqueduct-headless.default.svc.cluster.local"
	Port     string        // port suffix e.g. "4242"
	Interval time.Duration // polling interval
}

DiscoveryConfig holds DNS-based peer discovery parameters.

type PeerManager

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

PeerManager establishes and maintains outbound QUIC streams to peer addresses. Forwarding is zero-copy: the raw pooled []byte is written directly to peer streams with the MeshForwarded bit already set in buf[1]. Dynamic peer management uses RCU (Read-Copy-Update): Forward reads the atomic pointer; AddPeer/RemovePeer create a new slice and swap atomically.

func New

func New(ctx context.Context, addrs []string, tlsConf *tls.Config, quicConf *quic.Config) *PeerManager

New creates a PeerManager and begins background reconnect loops for each peer address.

func NewWithLogger added in v1.14.0

func NewWithLogger(ctx context.Context, addrs []string, tlsConf *tls.Config, quicConf *quic.Config, logger *slog.Logger) *PeerManager

NewWithLogger creates a PeerManager with a custom logger.

func (*PeerManager) ActivePeers

func (pm *PeerManager) ActivePeers() int

ActivePeers returns the number of currently connected peers.

func (*PeerManager) AddPeer added in v1.14.0

func (pm *PeerManager) AddPeer(ctx context.Context, addr string)

AddPeer dynamically adds a new peer address and starts its reconnect loop. Safe for concurrent use; uses RCU to swap the peer snapshot atomically.

func (*PeerManager) Close

func (pm *PeerManager) Close()

Close signals the reconnect goroutines to stop and drains the peer streams.

func (*PeerManager) Forward

func (pm *PeerManager) Forward(rawBuf []byte, addForwardedBit bool)

Forward sends rawBuf zero-copy to all connected peer streams. It sets the MeshForwarded bit in a copy of buf[1] to prevent mesh storms. RCU: reads the atomic pointer — no locks, no allocations on the hot path.

func (*PeerManager) ForwardRaw

func (pm *PeerManager) ForwardRaw(rawBuf []byte) error

ForwardRaw sends a pre-assembled raw byte slice to all connected peers, treating the MeshForwarded bit as already set in the caller's buffer. Used for in-test direct writes where the caller controls byte-level framing.

func (*PeerManager) PeerCount added in v1.14.0

func (pm *PeerManager) PeerCount() int

PeerCount returns the total number of managed peers (connected or not).

func (*PeerManager) RemovePeer added in v1.14.0

func (pm *PeerManager) RemovePeer(addr string)

RemovePeer dynamically removes a peer address, stops its reconnect loop, and closes its active stream. Safe for concurrent use via RCU.

type PeerRef

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

PeerRef holds the QUIC stream and connection state for a remote peer.

type Resolver added in v1.14.0

type Resolver interface {
	LookupHost(ctx context.Context, host string) ([]string, error)
}

Resolver abstracts DNS resolution for testability. The default implementation uses net.DefaultResolver (system DNS).

Jump to

Keyboard shortcuts

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