Files
jbergner 9030c2efb7
release-tag / release-image (push) Successful in 6m17s
1.5.3
2026-08-27 10:43:07 +02:00

468 lines
12 KiB
Go

package store
import (
"errors"
"fmt"
"math"
"sort"
"strings"
"time"
"neuroforge/internal/core"
"neuroforge/internal/vector"
)
var ErrDuplicateGoal = errors.New("duplicate goal")
func (s *Store) resolveConflictLocked(in *core.Memory) []core.Memory {
key := strings.TrimSpace(in.TruthKey)
if key == "" {
for _, tag := range in.Tags {
if strings.HasPrefix(strings.ToLower(tag), "truth:") && len(tag) > 6 {
key = strings.TrimSpace(tag[6:])
in.TruthKey = key
break
}
}
}
if key == "" {
return nil
}
var changed []core.Memory
group := "conflict:" + strings.ToLower(key)
for id, meta := range s.state.Memories {
if id == in.ID || !strings.EqualFold(strings.TrimSpace(meta.TruthKey), key) || meta.Status == core.MemoryArchived {
continue
}
old, ok := s.materializeMemoryLocked(id)
if !ok {
continue
}
if strings.EqualFold(strings.TrimSpace(old.Text), strings.TrimSpace(in.Text)) {
if in.Version >= old.Version {
old.Status = core.MemorySuperseded
in.Supersedes = appendUniqueString(in.Supersedes, old.ID)
changed = append(changed, cloneMemory(*old))
} else {
in.Status = core.MemorySuperseded
}
continue
}
if in.Version > old.Version {
old.Status = core.MemorySuperseded
in.Supersedes = appendUniqueString(in.Supersedes, old.ID)
changed = append(changed, cloneMemory(*old))
continue
}
if old.Version > in.Version {
in.Status = core.MemorySuperseded
continue
}
newScore := knowledgeScore(in)
oldScore := knowledgeScore(old)
if math.Abs(newScore-oldScore) < 0.15 {
old.Status = core.MemoryConflicted
old.ConflictGroup = group
in.Status = core.MemoryConflicted
in.ConflictGroup = group
changed = append(changed, cloneMemory(*old))
} else if newScore > oldScore {
old.Status = core.MemorySuperseded
in.Supersedes = appendUniqueString(in.Supersedes, old.ID)
changed = append(changed, cloneMemory(*old))
} else {
in.Status = core.MemorySuperseded
}
}
return changed
}
func knowledgeScore(m *core.Memory) float64 {
confidence := m.Confidence
if confidence <= 0 {
confidence = 1
}
salience := m.Salience
if salience <= 0 {
salience = 1
}
return 0.55*confidence + 0.25*vector.Clamp((m.Reward+1)/2, 0, 1) + 0.20*vector.Clamp(salience/2, 0, 1)
}
func appendUniqueString(xs []string, v string) []string {
for _, x := range xs {
if x == v {
return xs
}
}
return append(xs, v)
}
func (s *Store) ResolveConflict(truthKey, winnerID string) error {
s.mu.Lock()
defer s.mu.Unlock()
winner, ok := s.materializeMemoryLocked(winnerID)
if !ok || !strings.EqualFold(strings.TrimSpace(winner.TruthKey), strings.TrimSpace(truthKey)) {
return errors.New("winner memory not found for truth key")
}
changed := []core.Memory{}
for id, meta := range s.state.Memories {
if !strings.EqualFold(strings.TrimSpace(meta.TruthKey), strings.TrimSpace(truthKey)) {
continue
}
m, ok := s.materializeMemoryLocked(id)
if !ok {
continue
}
m.ConflictGroup = ""
if m.ID == winnerID {
m.Status = core.MemoryActive
} else {
m.Status = core.MemorySuperseded
winner.Supersedes = appendUniqueString(winner.Supersedes, m.ID)
}
changed = append(changed, cloneMemory(*m))
}
if len(changed) == 0 {
return errors.New("no memories found for truth key")
}
// Winner may have gained supersedes after its earlier copy was appended.
for i := range changed {
if changed[i].ID == winner.ID {
changed[i] = cloneMemory(*winner)
}
}
return s.commitLocked("memory.upsert", changed)
}
func (s *Store) ConflictsSnapshot() map[string][]core.Memory {
s.mu.RLock()
defer s.mu.RUnlock()
out := map[string][]core.Memory{}
for id, meta := range s.state.Memories {
if meta.Status != core.MemoryConflicted || meta.TruthKey == "" {
continue
}
if m, ok := s.fullMemoryForReadLocked(id); ok {
out[m.TruthKey] = append(out[m.TruthKey], cloneMemory(m))
}
}
return out
}
func (s *Store) SetMemoryHomeShard(id, shard string) error {
s.mu.Lock()
defer s.mu.Unlock()
m, ok := s.materializeMemoryLocked(id)
if !ok {
return errors.New("memory not found")
}
m.HomeShardID = shard
return s.commitLocked("memory.upsert", []core.Memory{cloneMemory(*m)})
}
func cloneGoal(g core.Goal) core.Goal {
g.MemoryIDs = append([]string(nil), g.MemoryIDs...)
g.Tags = append([]string(nil), g.Tags...)
return g
}
func (s *Store) UpsertGoal(g *core.Goal) error {
if g == nil || strings.TrimSpace(g.Title) == "" {
return errors.New("goal title required")
}
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UTC()
newGoal := g.ID == "" || s.state.Goals[g.ID] == nil
if newGoal {
fingerprint := goalFingerprint(g)
for id, existing := range s.state.Goals {
if existing == nil || existing.Status == core.GoalCompleted {
continue
}
if goalFingerprint(existing) == fingerprint {
return fmt.Errorf("%w: existing goal %s", ErrDuplicateGoal, id)
}
}
}
if g.ID == "" {
g.ID = NewID("goal")
}
if g.Status == "" {
g.Status = core.GoalActive
}
if g.Priority == 0 {
g.Priority = 50
}
if g.IntervalMinutes <= 0 {
g.IntervalMinutes = s.state.Config.Autonomy.DefaultGoalIntervalMinutes
}
if newGoal {
g.AutoRun = s.state.Config.Autonomy.RunOnGoalCreate
g.ResearchEnabled = s.state.Config.Research.Goal.Enabled
if g.AutoRun {
g.NextCycleAt = now
}
}
g.Progress = vector.Clamp(g.Progress, 0, 1)
if g.CreatedAt.IsZero() {
if old := s.state.Goals[g.ID]; old != nil {
g.CreatedAt = old.CreatedAt
} else {
g.CreatedAt = now
}
}
g.UpdatedAt = now
cp := cloneGoal(*g)
s.state.Goals[g.ID] = &cp
return s.commitLocked("goal.upsert", cp)
}
func goalFingerprint(g *core.Goal) string {
if g == nil {
return ""
}
normalize := func(v string) string {
return strings.ToLower(strings.Join(strings.Fields(strings.TrimSpace(v)), " "))
}
return normalize(g.Title) + "\x00" + normalize(g.Description) + "\x00" + normalize(g.Target)
}
func (s *Store) GetGoal(id string) (*core.Goal, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
g := s.state.Goals[id]
if g == nil {
return nil, false
}
cp := cloneGoal(*g)
return &cp, true
}
func (s *Store) GoalsSnapshot() []core.Goal {
s.mu.RLock()
defer s.mu.RUnlock()
out := make([]core.Goal, 0, len(s.state.Goals))
for _, g := range s.state.Goals {
out = append(out, cloneGoal(*g))
}
sort.Slice(out, func(i, j int) bool {
if out[i].Priority == out[j].Priority {
return out[i].UpdatedAt.After(out[j].UpdatedAt)
}
return out[i].Priority > out[j].Priority
})
return out
}
func (s *Store) PauseGoal(id string) (*core.Goal, error) {
s.mu.Lock()
defer s.mu.Unlock()
g := s.state.Goals[id]
if g == nil {
return nil, errors.New("goal not found")
}
if g.Status == core.GoalCompleted || g.Status == core.GoalFailed {
return nil, fmt.Errorf("goal in status %q cannot be paused", g.Status)
}
g.Status = core.GoalPaused
g.NextCycleAt = time.Time{}
g.UpdatedAt = time.Now().UTC()
cp := cloneGoal(*g)
if err := s.commitLocked("goal.upsert", cp); err != nil {
return nil, err
}
return &cp, nil
}
func (s *Store) ResumeGoal(id string) (*core.Goal, error) {
s.mu.Lock()
defer s.mu.Unlock()
g := s.state.Goals[id]
if g == nil {
return nil, errors.New("goal not found")
}
if g.Status == core.GoalCompleted || g.Status == core.GoalFailed {
return nil, fmt.Errorf("goal in status %q cannot be resumed", g.Status)
}
g.Status = core.GoalActive
g.ConsecutiveErrors = 0
g.LastError = ""
if g.AutoRun && s.state.Config.Autonomy.Enabled {
g.NextCycleAt = time.Now().UTC()
} else {
g.NextCycleAt = time.Time{}
}
g.UpdatedAt = time.Now().UTC()
cp := cloneGoal(*g)
if err := s.commitLocked("goal.upsert", cp); err != nil {
return nil, err
}
return &cp, nil
}
func (s *Store) DeleteGoal(id string) error {
s.mu.Lock()
defer s.mu.Unlock()
if _, ok := s.state.Goals[id]; !ok {
return errors.New("goal not found")
}
delete(s.state.Goals, id)
return s.commitLocked("goal.delete", id)
}
func (s *Store) AddLearningCycle(c core.LearningCycle) error {
s.mu.Lock()
defer s.mu.Unlock()
if c.ID == "" {
c.ID = NewID("cycle")
}
if c.CreatedAt.IsZero() {
c.CreatedAt = time.Now().UTC()
}
s.state.Cycles = append(s.state.Cycles, c)
if len(s.state.Cycles) > 10000 {
s.state.Cycles = s.state.Cycles[len(s.state.Cycles)-10000:]
}
return s.commitLocked("cycle.add", c)
}
func (s *Store) RecentLearningCycles(limit int) []core.LearningCycle {
s.mu.RLock()
defer s.mu.RUnlock()
if limit <= 0 || limit > len(s.state.Cycles) {
limit = len(s.state.Cycles)
}
out := append([]core.LearningCycle(nil), s.state.Cycles[len(s.state.Cycles)-limit:]...)
for i, j := 0, len(out)-1; i < j; i, j = i+1, j-1 {
out[i], out[j] = out[j], out[i]
}
return out
}
type RetentionResult struct {
Deleted int `json:"deleted"`
Compressed int `json:"compressed"`
IDs []string `json:"ids,omitempty"`
}
type retentionCandidate struct {
id string
utility float64
ageDays float64
}
func memoryUtility(m *core.Memory, now time.Time) float64 {
ageDays := now.Sub(m.AccessedAt).Hours() / 24
if ageDays < 0 {
ageDays = 0
}
freshness := math.Exp(-ageDays / 30)
access := math.Min(1, math.Log1p(float64(m.AccessCount))/4)
confidence := m.Confidence
if confidence <= 0 {
confidence = 1
}
reward := (m.Reward + 1) / 2
return 0.30*freshness + 0.20*access + 0.20*vector.Clamp(m.Salience/2, 0, 1) + 0.20*vector.Clamp(confidence, 0, 1) + 0.10*vector.Clamp(reward, 0, 1)
}
func (s *Store) RunRetention(now time.Time) (RetentionResult, error) {
s.mu.Lock()
defer s.mu.Unlock()
cfg := s.state.Config.Retention
result := RetentionResult{}
if !cfg.Enabled {
return result, nil
}
if now.IsZero() {
now = time.Now().UTC()
}
candidates := make([]retentionCandidate, 0, len(s.state.Memories))
deleteSet := map[string]bool{}
changed := []core.Memory{}
for id, m := range s.state.Memories {
if m == nil {
continue
}
ageDays := now.Sub(m.CreatedAt).Hours() / 24
if m.MemoryType == core.MemoryWorking && cfg.WorkingTTLHours > 0 && now.Sub(m.CreatedAt).Hours() >= cfg.WorkingTTLHours {
deleteSet[id] = true
continue
}
utility := memoryUtility(m, now)
candidates = append(candidates, retentionCandidate{id: id, utility: utility, ageDays: ageDays})
if ageDays < cfg.MinAgeDays || utility >= cfg.MinUtility || m.Status == core.MemoryArchived {
continue
}
if cfg.DeleteConsolidated && m.ConsolidatedInto != "" {
deleteSet[id] = true
continue
}
full, ok := s.materializeMemoryLocked(m.ID)
if !ok {
continue
}
m = full
m.Status = core.MemoryArchived
m.Compressed = true
m.Vector = nil
m.VectorDim = 0
if cfg.CompressChars > 0 && len([]rune(m.Text)) > cfg.CompressChars {
r := []rune(m.Text)
m.Text = string(r[:cfg.CompressChars]) + "…"
}
s.trackHotMemoryLocked(m.ID, m)
changed = append(changed, cloneMemory(*m))
result.Compressed++
result.IDs = append(result.IDs, m.ID)
}
if cfg.MaxMemories > 0 && len(s.state.Memories)-len(deleteSet) > cfg.MaxMemories {
sort.Slice(candidates, func(i, j int) bool { return candidates[i].utility < candidates[j].utility })
need := len(s.state.Memories) - len(deleteSet) - cfg.MaxMemories
for _, c := range candidates {
if need <= 0 {
break
}
m := s.state.Memories[c.id]
if m == nil || deleteSet[c.id] || m.MemoryType == core.MemoryProcedural {
continue
}
deleteSet[c.id] = true
need--
}
}
if len(changed) > 0 {
if err := s.commitLocked("memory.upsert", changed); err != nil {
return result, err
}
}
if len(deleteSet) > 0 {
ids := make([]string, 0, len(deleteSet))
for id := range deleteSet {
ids = append(ids, id)
delete(s.state.Memories, id)
s.untrackHotMemoryLocked(id)
if s.pageCache != nil {
s.pageCache.Delete(id)
}
for k, syn := range s.state.Synapses {
if syn.A == id || syn.B == id {
delete(s.state.Synapses, k)
}
}
}
sort.Strings(ids)
result.Deleted = len(ids)
result.IDs = append(result.IDs, ids...)
if err := s.commitLocked("memory.delete", ids); err != nil {
return result, err
}
}
if len(changed) > 0 || len(deleteSet) > 0 {
s.rebuildIndexesLocked()
}
return result, nil
}