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

294 lines
8.2 KiB
Go

package store
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"time"
"neuroforge/internal/core"
)
type ClusterDecision struct {
EntryID string `json:"entry_id"`
Term uint64 `json:"term"`
Index uint64 `json:"index"`
Decision string `json:"decision"`
CreatedAt time.Time `json:"created_at"`
}
func (s *Store) clusterDir() string { return filepath.Join(s.dir, "cluster") }
func (s *Store) pendingClusterDir() string { return filepath.Join(s.clusterDir(), "pending") }
func (s *Store) decisionClusterDir() string { return filepath.Join(s.clusterDir(), "decisions") }
func writeJSONSync(path string, v any) error {
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return err
}
b, err := json.MarshalIndent(v, "", " ")
if err != nil {
return err
}
tmp := path + ".tmp"
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0600)
if err != nil {
return err
}
if _, err := f.Write(b); err == nil {
err = f.Sync()
}
closeErr := f.Close()
if err != nil {
return err
}
if closeErr != nil {
return closeErr
}
if err := os.Rename(tmp, path); err != nil {
return err
}
if d, err := os.Open(filepath.Dir(path)); err == nil {
_ = d.Sync()
_ = d.Close()
}
return nil
}
func (s *Store) ClusterState() core.ClusterState {
s.mu.RLock()
defer s.mu.RUnlock()
return s.state.Cluster
}
func (s *Store) NextClusterIndex(term uint64) (uint64, error) {
s.mu.RLock()
defer s.mu.RUnlock()
if term < s.state.Cluster.Term {
return 0, fmt.Errorf("stale cluster term %d < %d", term, s.state.Cluster.Term)
}
last := s.state.Cluster.LastIndex
if s.state.Cluster.CommitIndex > last {
last = s.state.Cluster.CommitIndex
}
return last + 1, nil
}
func (s *Store) PrepareClusterEntry(entry core.ClusterEntry) error {
if entry.ID == "" || entry.Type == "" || entry.Term == 0 || entry.Index == 0 || entry.LeaderID == "" {
return errors.New("invalid cluster entry")
}
s.mu.RLock()
cfg := s.state.Config.Cluster
state := s.state.Cluster
s.mu.RUnlock()
if !cfg.Enabled {
return errors.New("cluster is disabled")
}
if entry.Term < state.Term {
return fmt.Errorf("stale cluster term %d < %d", entry.Term, state.Term)
}
leaderID := cfg.LeaderID
if cfg.AutoElection && state.LeaderID != "" {
leaderID = state.LeaderID
}
if entry.LeaderID != leaderID {
return fmt.Errorf("entry leader %q does not match current leader %q", entry.LeaderID, leaderID)
}
if entry.Index <= state.CommitIndex {
// Idempotent retry of an already committed entry is acceptable only if
// the decision exists locally.
if d, ok := s.ClusterDecision(entry.ID); ok && d.Decision == "commit" {
return nil
}
return fmt.Errorf("cluster index %d is already committed through %d", entry.Index, state.CommitIndex)
}
if err := s.appendClusterLogEntry(entry); err != nil {
return err
}
return writeJSONSync(filepath.Join(s.pendingClusterDir(), entry.ID+".json"), &entry)
}
func (s *Store) AbortPreparedClusterEntry(id string) error {
if strings.TrimSpace(id) == "" {
return errors.New("entry id required")
}
err := os.Remove(filepath.Join(s.pendingClusterDir(), id+".json"))
if errors.Is(err, os.ErrNotExist) {
return nil
}
return err
}
func (s *Store) PendingClusterEntries() []core.ClusterEntry {
ents, err := os.ReadDir(s.pendingClusterDir())
if err != nil {
return nil
}
out := make([]core.ClusterEntry, 0, len(ents))
for _, ent := range ents {
if ent.IsDir() || !strings.HasSuffix(ent.Name(), ".json") {
continue
}
var e core.ClusterEntry
if s.loadJSON(filepath.Join(s.pendingClusterDir(), ent.Name()), &e) == nil && e.ID != "" {
out = append(out, e)
}
}
sort.Slice(out, func(i, j int) bool {
if out[i].Term == out[j].Term {
return out[i].Index < out[j].Index
}
return out[i].Term < out[j].Term
})
return out
}
func (s *Store) RecordClusterDecision(entry core.ClusterEntry, decision string) error {
if decision != "commit" && decision != "abort" {
return errors.New("cluster decision must be commit or abort")
}
d := ClusterDecision{EntryID: entry.ID, Term: entry.Term, Index: entry.Index, Decision: decision, CreatedAt: time.Now().UTC()}
if err := s.appendClusterLogDecision(d); err != nil {
return err
}
return writeJSONSync(filepath.Join(s.decisionClusterDir(), entry.ID+".json"), &d)
}
func (s *Store) ClusterDecision(id string) (ClusterDecision, bool) {
var d ClusterDecision
if err := s.loadJSON(filepath.Join(s.decisionClusterDir(), id+".json"), &d); err != nil {
return ClusterDecision{}, false
}
return d, d.EntryID != ""
}
func (s *Store) CommitPreparedClusterEntry(entry core.ClusterEntry) error {
if d, ok := s.ClusterDecision(entry.ID); ok && d.Decision == "abort" {
return errors.New("cluster entry was aborted")
}
var pending core.ClusterEntry
pendingPath := filepath.Join(s.pendingClusterDir(), entry.ID+".json")
if err := s.loadJSON(pendingPath, &pending); err != nil {
// Leader may recover after applying the memory but before deleting pending.
state := s.ClusterState()
if state.CommitIndex >= entry.Index {
return nil
}
return fmt.Errorf("prepared cluster entry missing: %w", err)
}
if pending.Term != entry.Term || pending.Index != entry.Index || pending.Type != entry.Type || pending.LeaderID != entry.LeaderID {
return errors.New("prepared cluster entry does not match commit")
}
switch entry.Type {
case "memory.upsert":
var m core.Memory
if err := json.Unmarshal(entry.Payload, &m); err != nil {
return err
}
if err := s.UpsertClusterMemory(&m); err != nil {
return err
}
default:
return fmt.Errorf("unsupported cluster entry type %q", entry.Type)
}
s.mu.Lock()
if entry.Term > s.state.Cluster.Term {
s.state.Cluster.Term = entry.Term
}
if entry.Index > s.state.Cluster.LastIndex {
s.state.Cluster.LastIndex = entry.Index
}
if entry.Index > s.state.Cluster.CommitIndex {
s.state.Cluster.CommitIndex = entry.Index
}
s.state.Cluster.LastCommit = time.Now().UTC()
state := s.state.Cluster
err := s.commitLocked("cluster.state", state)
s.mu.Unlock()
if err != nil {
return err
}
_ = os.Remove(pendingPath)
return nil
}
func (s *Store) UpsertClusterMemory(m *core.Memory) error {
if m == nil || strings.TrimSpace(m.ID) == "" {
return errors.New("cluster memory id required")
}
if existing, ok := s.GetMemory(m.ID); ok {
// A commit retry is idempotent only when the immutable learning payload
// matches. Location/timestamps may legitimately differ between replicas.
if sameClusterMemory(existing, m) {
return nil
}
return fmt.Errorf("cluster memory %s already exists with different content", m.ID)
}
cp := cloneMemory(*m)
cfg := s.Config()
if cp.OriginShardID == "" {
cp.OriginShardID = cp.ShardID
}
if cp.OriginShardID == "" {
cp.OriginShardID = s.EffectiveLeaderID()
}
cp.ShardID = cfg.Sharding.LocalShardID
if cp.HomeShardID == "" {
cp.HomeShardID = s.EffectiveLeaderID()
}
return s.AddMemory(&cp)
}
func sameClusterMemory(a, b *core.Memory) bool {
if a == nil || b == nil {
return a == b
}
if a.ID != b.ID || a.Text != b.Text || a.Kind != b.Kind || a.MemoryType != b.MemoryType ||
a.TruthKey != b.TruthKey || a.Version != b.Version || len(a.Vector) != len(b.Vector) {
return false
}
for i := range a.Vector {
if a.Vector[i] != b.Vector[i] {
return false
}
}
return true
}
func (s *Store) ClusterStatus() map[string]any {
s.mu.RLock()
cfg := s.state.Config.Cluster
state := s.state.Cluster
s.mu.RUnlock()
voters := 1
peers := 0
for _, p := range cfg.Peers {
if !p.Enabled {
continue
}
peers++
if p.Voting {
voters++
}
}
quorum := cfg.Quorum
if quorum <= 0 {
quorum = voters/2 + 1
}
leaderID := cfg.LeaderID
if cfg.AutoElection {
leaderID = state.LeaderID
}
return map[string]any{
"enabled": cfg.Enabled, "node_id": cfg.NodeID, "leader_id": leaderID, "configured_leader_id": cfg.LeaderID, "auto_election": cfg.AutoElection, "role": state.Role, "configured_term": cfg.Term,
"term": state.Term, "voted_for": state.VotedFor, "last_heartbeat": state.LastHeartbeat, "last_index": state.LastIndex, "commit_index": state.CommitIndex, "last_commit": state.LastCommit,
"peers": peers, "voters": voters, "quorum": quorum, "pending": len(s.PendingClusterEntries()), "replicated_log": s.ClusterLogStats(),
}
}