Files
jbergner a6bc71fb3a
release-tag / Resolve release metadata (push) Successful in 30s
release-tag / Build knowledge (push) Failing after 4m51s
release-tag / Build control (push) Failing after 5m0s
release-tag / Build agent (push) Failing after 5m0s
release-tag / Build agent-data-init (push) Failing after 5m5s
release-tag / Build neuroforge-worker (push) Failing after 5m7s
release-tag / Build neuroforge (push) Failing after 5m9s
Init
2026-08-26 18:34:41 +02:00

196 lines
5.3 KiB
Go

package brain
import (
"context"
"fmt"
"hash/fnv"
"sync"
"time"
"neuroforge/internal/core"
"neuroforge/internal/store"
)
func electionTimeout(cfg core.Config, term uint64) time.Duration {
minMS, maxMS := cfg.Cluster.ElectionMinMS, cfg.Cluster.ElectionMaxMS
if minMS <= 0 {
minMS = 1200
}
if maxMS <= minMS {
maxMS = minMS * 2
}
h := fnv.New64a()
_, _ = h.Write([]byte(fmt.Sprintf("%s:%d:%d", cfg.Cluster.NodeID, term, time.Now().UnixNano()/int64(time.Millisecond))))
span := uint64(maxMS - minMS)
jitter := 0
if span > 0 {
jitter = int(h.Sum64() % span)
}
return time.Duration(minMS+jitter) * time.Millisecond
}
func (e *Engine) resetElectionDeadline(cfg core.Config, state core.ClusterState, now time.Time) {
e.electionMu.Lock()
defer e.electionMu.Unlock()
if e.electionDeadline.IsZero() || state.LastHeartbeat.After(e.observedHeartbeat) {
e.observedHeartbeat = state.LastHeartbeat
e.electionDeadline = now.Add(electionTimeout(cfg, state.Term))
}
}
func (e *Engine) electionDue(now time.Time) bool {
e.electionMu.Lock()
defer e.electionMu.Unlock()
if e.electionRunning || e.electionDeadline.IsZero() || now.Before(e.electionDeadline) {
return false
}
e.electionRunning = true
return true
}
func (e *Engine) electionFinished(cfg core.Config, term uint64) {
e.electionMu.Lock()
defer e.electionMu.Unlock()
e.electionRunning = false
e.electionDeadline = time.Now().Add(electionTimeout(cfg, term))
}
func (e *Engine) attemptElection(ctx context.Context) {
cfg := e.store.Config()
req, err := e.store.StartElection()
if err != nil {
e.electionFinished(cfg, e.store.ClusterState().Term)
return
}
defer e.electionFinished(cfg, req.Term)
voters, quorum := clusterVoters(cfg)
votes := 1
var mu sync.Mutex
var wg sync.WaitGroup
for _, peer := range cfg.Cluster.Peers {
if !peer.Enabled || !peer.Voting {
continue
}
wg.Add(1)
go func(p core.ClusterPeer) {
defer wg.Done()
var resp core.ClusterVoteResponse
if err := e.clusterPost(ctx, p.BaseURL, "/internal/v1/cluster/request-vote", req, &resp); err != nil {
return
}
if resp.Term > req.Term {
_ = e.store.StepDown(resp.Term, "")
return
}
if resp.Term == req.Term && resp.VoteGranted {
mu.Lock()
votes++
mu.Unlock()
}
}(peer)
}
wg.Wait()
if e.store.ClusterState().Term != req.Term {
return
}
if votes >= quorum {
if err := e.store.BecomeLeader(req.Term); err == nil {
e.sendHeartbeats(ctx)
}
return
}
_ = voters // retained in status/debugging; quorum already derived from same set.
}
func (e *Engine) sendHeartbeats(ctx context.Context) {
cfg := e.store.Config()
state := e.store.ClusterState()
if !cfg.Cluster.Enabled || !cfg.Cluster.AutoElection || state.Role != store.ClusterLeader || state.LeaderID != cfg.Cluster.NodeID {
return
}
h := core.ClusterHeartbeat{Term: state.Term, LeaderID: cfg.Cluster.NodeID, CommitIndex: state.CommitIndex, LastIndex: state.LastIndex}
var wg sync.WaitGroup
for _, peer := range cfg.Cluster.Peers {
if !peer.Enabled {
continue
}
wg.Add(1)
go func(p core.ClusterPeer) {
defer wg.Done()
var resp core.ClusterHeartbeatResponse
if err := e.clusterPost(ctx, p.BaseURL, "/internal/v1/cluster/heartbeat", h, &resp); err != nil {
return
}
if resp.Term > h.Term {
_ = e.store.StepDown(resp.Term, "")
}
}(peer)
}
wg.Wait()
_ = e.store.TouchLeaderHeartbeat()
}
func (e *Engine) ClusterVote(req core.ClusterVoteRequest) (core.ClusterVoteResponse, error) {
return e.store.GrantVote(req)
}
func (e *Engine) ClusterHeartbeat(h core.ClusterHeartbeat) (core.ClusterHeartbeatResponse, error) {
return e.store.AcceptHeartbeat(h)
}
func (e *Engine) RunV5Maintenance(ctx context.Context) {
go e.RunV4Maintenance(ctx)
ticker := time.NewTicker(100 * time.Millisecond)
defer ticker.Stop()
var lastHeartbeatSent, lastTier, lastIndexMerge time.Time
for {
select {
case <-ctx.Done():
return
case now := <-ticker.C:
cfg := e.store.Config()
if cfg.Cluster.Enabled && cfg.Cluster.AutoElection {
state := e.store.ClusterState()
if state.Role == store.ClusterLeader && state.LeaderID == cfg.Cluster.NodeID {
hb := time.Duration(cfg.Cluster.HeartbeatMS) * time.Millisecond
if hb <= 0 {
hb = 350 * time.Millisecond
}
if now.Sub(lastHeartbeatSent) >= hb {
e.sendHeartbeats(ctx)
lastHeartbeatSent = now
}
} else {
e.resetElectionDeadline(cfg, state, now)
if e.electionDue(now) {
go e.attemptElection(ctx)
}
}
}
if cfg.Storage.Tiering.Enabled {
iv := time.Duration(cfg.Storage.Tiering.IntervalMinutes) * time.Minute
if iv <= 0 {
iv = 5 * time.Minute
}
if lastTier.IsZero() || now.Sub(lastTier) >= iv {
e.store.TierMemoryBodies(now)
lastTier = now
}
}
if cfg.Storage.IndexSegments.Enabled && cfg.Storage.IndexSegments.BackgroundMergeMinutes > 0 {
iv := time.Duration(cfg.Storage.IndexSegments.BackgroundMergeMinutes) * time.Minute
if lastIndexMerge.IsZero() || now.Sub(lastIndexMerge) >= iv {
st := e.store.IndexSnapshotStatus()
deltas, _ := st["deltas"].(int)
threshold := cfg.Storage.IndexSegments.MergeAtDeltas
if threshold <= 0 {
threshold = 8
}
if deltas >= threshold {
_, _ = e.store.CompactIndexSegments()
}
lastIndexMerge = now
}
}
}
}
}