Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ require (
github.com/hashicorp/go-multierror v1.1.1
github.com/hashicorp/serf v0.10.2
github.com/opencontainers/runtime-spec v1.2.1
github.com/prometheus/client_golang v1.20.2
github.com/spf13/viper v1.20.0
github.com/vistara-labs/firecracker-containerd v0.0.0-20240707190021-1287a7cb7490
google.golang.org/grpc v1.71.0
Expand All @@ -28,6 +29,8 @@ require (
github.com/Microsoft/hcsshim v0.12.9 // indirect
github.com/armon/go-metrics v0.4.1 // indirect
github.com/asaskevich/govalidator v0.0.0-20230301143203-a9d515a09cc2 // indirect
github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/containerd/cgroups/v3 v3.0.5 // indirect
github.com/containerd/console v1.0.4 // indirect
github.com/containerd/continuity v0.4.5 // indirect
Expand Down Expand Up @@ -84,13 +87,17 @@ require (
github.com/moby/sys/userns v0.1.0 // indirect
github.com/moby/term v0.0.0-20221205130635-1aeaba878587 // indirect
github.com/morikuni/aec v1.0.0 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/opencontainers/go-digest v1.0.0 // indirect
github.com/opencontainers/image-spec v1.1.1 // indirect
github.com/opencontainers/runc v1.2.6 // indirect
github.com/opencontainers/selinux v1.12.0 // indirect
github.com/opentracing/opentracing-go v1.2.0 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
github.com/prometheus/client_model v0.6.1 // indirect
github.com/prometheus/common v0.55.0 // indirect
github.com/prometheus/procfs v0.15.1 // indirect
github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529 // indirect
github.com/shirou/gopsutil v3.21.11+incompatible // indirect
github.com/stretchr/testify v1.10.0 // indirect
Expand Down
12 changes: 12 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ github.com/aws/aws-sdk-go v1.15.11/go.mod h1:mFuSZ37Z9YOHbQEwBWztmVzqXrEkub65tZo
github.com/beorn7/perks v0.0.0-20160804104726-4c0e84591b9a/go.mod h1:Dwedo/Wpr24TaqPxmxbtue+5NUziq4I4S80YR8gNf3Q=
github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973/go.mod h1:Dwedo/Wpr24TaqPxmxbtue+5NUziq4I4S80YR8gNf3Q=
github.com/beorn7/perks v1.0.0/go.mod h1:KWe93zE9D1o94FZ5RNwFwVgaQK1VOXiVxmqh+CedLV8=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/bgentry/speakeasy v0.1.0/go.mod h1:+zsyZBPWlz7T6j88CTgSN5bM796AkVf0kBD4zp0CCIs=
github.com/bitly/go-simplejson v0.5.0/go.mod h1:cXHtHw4XUPsvGaxgjIAn8PhEWG9NfngEKAMDJEczWVA=
Expand All @@ -106,6 +107,8 @@ github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyY
github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU=
github.com/cespare/xxhash v1.1.0/go.mod h1:XrSqR1VqqWfGrhpAt58auRo0WTKS1nRRg3ghfAqPWnc=
github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/checkpoint-restore/go-criu/v4 v4.1.0/go.mod h1:xUQBLp4RLc5zJtWY++yjOoMoB5lihDt7fai+75m+rGw=
github.com/chzyer/logex v1.1.10/go.mod h1:+Ywpsq7O8HXn0nuIou7OrIPyXbp3wmkHB+jjWRnGsAI=
github.com/chzyer/readline v0.0.0-20180603132655-2972be24d48e/go.mod h1:nSuG5e5PlCu98SY8svDHJxuZscDgtXS6KTTbou5AhLI=
Expand Down Expand Up @@ -596,6 +599,8 @@ github.com/kr/pty v1.1.5/go.mod h1:9r2w37qlBe7rQ6e1fg1S/9xpWHSnaqNdHD3WcMdbPDA=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
github.com/magiconair/properties v1.8.0/go.mod h1:PppfXfuXeibc/6YijjN8zIbojt8czPbwD3XqdrwzmxQ=
github.com/mailru/easyjson v0.0.0-20190614124828-94de47d64c63/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc=
github.com/mailru/easyjson v0.0.0-20190626092158-b2ccc519800e/go.mod h1:C1wdFJiN94OJF2b5HbByQZoLdCWB1Yqtg26g4irojpc=
Expand Down Expand Up @@ -663,6 +668,7 @@ github.com/morikuni/aec v1.0.0 h1:nP9CBfwrvYnBRgY6qfDQkygYDmYwOilePFkwzv4dU8A=
github.com/morikuni/aec v1.0.0/go.mod h1:BbKIizmSmc5MMPqRYbxO4ZU0S0+P200+tUnFx7PXmsc=
github.com/mrunalp/fileutils v0.5.0/go.mod h1:M1WthSahJixYnrXQl/DFQuteStB1weuxD2QJNHXfbSQ=
github.com/munnerz/goautoneg v0.0.0-20120707110453-a547fc61f48d/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U=
github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U=
Expand Down Expand Up @@ -751,11 +757,15 @@ github.com/prometheus/client_golang v1.1.0/go.mod h1:I1FGZT9+L76gKKOs5djB6ezCbFQ
github.com/prometheus/client_golang v1.4.0/go.mod h1:e9GMxYsXl05ICDXkRhurwBS4Q3OK1iX/F2sw+iXX5zU=
github.com/prometheus/client_golang v1.7.1/go.mod h1:PY5Wy2awLA44sXw4AOSfFBetzPP4j5+D6mVACh+pe2M=
github.com/prometheus/client_golang v1.11.1/go.mod h1:Z6t4BnS23TR94PD6BsDNk8yVqroYurpAkEiz0P2BEV0=
github.com/prometheus/client_golang v1.20.2 h1:5ctymQzZlyOON1666svgwn3s6IKWgfbjsejTMiXIyjg=
github.com/prometheus/client_golang v1.20.2/go.mod h1:PIEt8X02hGcP8JWbeHyeZ53Y/jReSnHgO035n//V5WE=
github.com/prometheus/client_model v0.0.0-20171117100541-99fa1f4be8e5/go.mod h1:MbSGuTsp3dbXC40dX6PRTWyKYBIrTGTE9sqQNg2J8bo=
github.com/prometheus/client_model v0.0.0-20180712105110-5c3871d89910/go.mod h1:MbSGuTsp3dbXC40dX6PRTWyKYBIrTGTE9sqQNg2J8bo=
github.com/prometheus/client_model v0.0.0-20190129233127-fd36f4220a90/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA=
github.com/prometheus/client_model v0.0.0-20190812154241-14fe0d1b01d4/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA=
github.com/prometheus/client_model v0.2.0/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA=
github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E=
github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY=
github.com/prometheus/common v0.0.0-20180110214958-89604d197083/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro=
github.com/prometheus/common v0.0.0-20181113130724-41aa239b4cce/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro=
github.com/prometheus/common v0.4.0/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4=
Expand All @@ -764,6 +774,8 @@ github.com/prometheus/common v0.6.0/go.mod h1:eBmuwkDJBwy6iBfxCBob6t6dR6ENT/y+J+
github.com/prometheus/common v0.9.1/go.mod h1:yhUN8i9wzaXS3w1O07YhxHEBxD+W35wd8bs7vj7HSQ4=
github.com/prometheus/common v0.10.0/go.mod h1:Tlit/dnDKsSWFlCLTWaA1cyBgKHSMdTB80sz/V91rCo=
github.com/prometheus/common v0.26.0/go.mod h1:M7rCNAaPfAosfx8veZJCuw84e35h3Cfd9VFqTh1DIvc=
github.com/prometheus/common v0.55.0 h1:KEi6DK7lXW/m7Ig5i47x0vRzuBsHuvJdi5ee6Y3G1dc=
github.com/prometheus/common v0.55.0/go.mod h1:2SECS4xJG1kd8XF9IcM1gMX6510RAEL65zxzNImwdc8=
github.com/prometheus/procfs v0.0.0-20180125133057-cb4147076ac7/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk=
github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk=
github.com/prometheus/procfs v0.0.0-20190507164030-5867b95ac084/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA=
Expand Down
131 changes: 117 additions & 14 deletions pkg/cluster/serf.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ package cluster

import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
Expand All @@ -10,6 +12,7 @@ import (
"net"
"net/http"
"runtime"
"sort"
"strconv"
"strings"
"sync"
Expand All @@ -21,6 +24,7 @@ import (
"github.com/containerd/containerd/cio"
"github.com/google/uuid"
"github.com/hashicorp/serf/serf"
"github.com/prometheus/client_golang/prometheus"
log "github.com/sirupsen/logrus"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/anypb"
Expand All @@ -31,7 +35,8 @@ const (
SpawnRequestLabel = "hypercore-request-payload"
StateBroadcastEvent = "hypercore_state_broadcast"

WorkloadBroadcastPeriod = time.Second * 5
WorkloadBroadcastPeriod = time.Second * 30 // Increased from 5s to 30s
MaxQueueDepth = 3600 // Increased for production safety
)

type SavedStatusUpdate struct {
Expand All @@ -56,6 +61,31 @@ type Agent struct {
lastStateSelf *pb.NodeStateResponse
lastStateUpdate map[string]SavedStatusUpdate
tmpStateUpdates map[string]*pb.NodeStateResponse
lastStateHash string // Track state hash to detect changes
stateMu sync.Mutex

// Prometheus metrics
serfQueueDepth prometheus.Gauge
workloadCount prometheus.Gauge
broadcastSkipped prometheus.Counter
stateChanges prometheus.Counter
}

// hashWorkloadState creates a consistent hash of the workload state
func (a *Agent) hashWorkloadState(state *pb.NodeStateResponse) string {
// Create a deterministic string representation
workloadIDs := make([]string, 0, len(state.GetWorkloads()))
for _, workload := range state.GetWorkloads() {
workloadIDs = append(workloadIDs, workload.GetId())
}

// Sort for consistency
sort.Strings(workloadIDs)

// Create hash
hash := sha256.New()
hash.Write([]byte(fmt.Sprintf("%d:%s", len(workloadIDs), strings.Join(workloadIDs, ","))))
return hex.EncodeToString(hash.Sum(nil))
}

func NewAgent(logger *log.Logger, baseURL, bindAddr string, respawn bool, repo *vcontainerd.Repo, tlsConfig *TLSConfig) (*Agent, error) {
Expand All @@ -82,27 +112,57 @@ func NewAgent(logger *log.Logger, baseURL, bindAddr string, respawn bool, repo *
cfg.MemberlistConfig.BindAddr = addr
cfg.MemberlistConfig.BindPort = bindPort
cfg.MemberlistConfig.AdvertisePort = bindPort
cfg.UserEventSizeLimit = 9 * 1024
cfg.UserEventSizeLimit = 2048 // Increased event size limit

// Conservative Serf configuration to reduce queue buildup
cfg.MemberlistConfig.GossipInterval = time.Second * 2 // Conservative gossip interval
cfg.MemberlistConfig.ProbeInterval = time.Second * 5 // Conservative probe interval
cfg.MemberlistConfig.SuspicionMult = 6 // Increased suspicion multiplier for stability
cfg.MemberlistConfig.GossipNodes = 2 // Reduce gossip nodes to decrease load

cfg.Init()

serf, err := serf.Create(cfg)
if err != nil {
return nil, err
}

// Initialize Prometheus metrics
serfQueueDepth := prometheus.NewGauge(prometheus.GaugeOpts{
Name: "hypercore_serf_queue_depth",
Help: "Current Serf event queue depth",
})
workloadCount := prometheus.NewGauge(prometheus.GaugeOpts{
Name: "hypercore_workload_count",
Help: "Number of running workloads on this node",
})
broadcastSkipped := prometheus.NewCounter(prometheus.CounterOpts{
Name: "hypercore_broadcast_skipped_total",
Help: "Total number of broadcasts skipped due to queue depth",
})
stateChanges := prometheus.NewCounter(prometheus.CounterOpts{
Name: "hypercore_state_changes_total",
Help: "Total number of state changes detected",
})

// Register metrics
prometheus.MustRegister(serfQueueDepth, workloadCount, broadcastSkipped, stateChanges)

agent := &Agent{
eventCh: eventCh,
cfg: cfg,
baseURL: baseURL,
serviceProxy: serviceProxy,
serf: serf,
logger: logger,
ctrRepo: repo,
lastStateUpdate: make(map[string]SavedStatusUpdate),
tmpStateUpdates: make(map[string]*pb.NodeStateResponse),
}
go agent.monitorWorkloads()
go agent.monitorStateUpdates(respawn)
eventCh: eventCh,
cfg: cfg,
baseURL: baseURL,
serviceProxy: serviceProxy,
serf: serf,
logger: logger,
ctrRepo: repo,
lastStateUpdate: make(map[string]SavedStatusUpdate),
tmpStateUpdates: make(map[string]*pb.NodeStateResponse),
serfQueueDepth: serfQueueDepth,
workloadCount: workloadCount,
broadcastSkipped: broadcastSkipped,
stateChanges: stateChanges,
}

return agent, nil
}
Expand Down Expand Up @@ -669,6 +729,49 @@ func (a *Agent) monitorWorkloads() {
a.lastStateSelf = &resp
a.lastStateMu.Unlock()

// Check if state has actually changed
currentHash := a.hashWorkloadState(&resp)
a.stateMu.Lock()
stateChanged := currentHash != a.lastStateHash
if stateChanged {
a.lastStateHash = currentHash
}
a.stateMu.Unlock()

// Only broadcast if state changed
if !stateChanged {
a.logger.Debug("State unchanged, skipping broadcast")
continue
}

// Check queue depth before broadcasting
stats := a.serf.Stats()
if queueDepthStr, ok := stats["event_queue_depth"]; ok {
if queueDepth, err := strconv.Atoi(queueDepthStr); err == nil {
// Update Prometheus metric
a.serfQueueDepth.Set(float64(queueDepth))

// Log queue depth for monitoring
if queueDepth > 1000 {
a.logger.Infof("Queue depth: %d (monitoring)", queueDepth)
}

if queueDepth > MaxQueueDepth {
a.logger.Warnf("Queue depth %d exceeds limit %d, skipping broadcast to prevent message drops", queueDepth, MaxQueueDepth)
a.broadcastSkipped.Inc()
continue
}
}
}

// Update workload count metric
a.workloadCount.Set(float64(len(resp.GetWorkloads())))

// Increment state changes counter
if stateChanged {
a.stateChanges.Inc()
}

// batch size of 10
parts := int(math.Ceil(float64(len(resp.GetWorkloads())) / 10))

Expand Down
Loading