Documentation
¶
Index ¶
- type DefaultResolver
- type Discovery
- type DiscoveryConfig
- type PeerManager
- func (pm *PeerManager) ActivePeers() int
- func (pm *PeerManager) AddPeer(ctx context.Context, addr string)
- func (pm *PeerManager) Close()
- func (pm *PeerManager) Forward(rawBuf []byte, addForwardedBit bool)
- func (pm *PeerManager) ForwardRaw(rawBuf []byte) error
- func (pm *PeerManager) PeerCount() int
- func (pm *PeerManager) RemovePeer(addr string)
- type PeerRef
- type Resolver
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
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:
- Every N seconds, resolve host → []IP via net.LookupHost.
- Sort + deduplicate the result set.
- Diff against the currently known set.
- For IPs not in current set → PeerManager.AddPeer.
- 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
KnownPeers returns a sorted copy of the currently discovered peer addresses.
func (*Discovery) ResolveOnce ¶ added in v1.14.0
ResolveOnce performs a single resolution + reconciliation (exported for testing).
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.