Files
jbergner 440423c5b6
All checks were successful
release-tag / release-image (push) Successful in 2m43s
RC-3
2026-08-09 11:29:13 +02:00

1098 lines
29 KiB
Go

package graph
import (
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"math"
"sort"
"strings"
"sync"
"time"
"github.com/local/glpi-neural-brain/internal/model"
)
type Store struct {
mu sync.RWMutex
nodes map[string]model.Node
edges map[string]model.Edge
vectors map[string][]float32
version uint64
persistedVersion uint64
pairCursor int
embeddingModel string
embeddingDigest string
embeddingMetaGeneration uint64
mutations MutationStats
db *sql.DB
analysisDB *sql.DB
dbPath string
journalMode string
dirtyNodes map[string]uint64
dirtyEdges map[string]uint64
dirtyVectors map[string]uint64
deletedNodes map[string]uint64
deletedEdges map[string]uint64
deletedVectors map[string]uint64
analysisMu sync.Mutex
analysisQueue chan analysisRecord
analysisWG sync.WaitGroup
analysisLastMutations MutationStats
analysisPendingChanges []GraphChange
analysisPendingTruncated int
analysisChangesDropped uint64
analysisDropped uint64
analysisQueueDropped uint64
analysisPersistDropped uint64
analysisLastDropReason string
analysisLastDropAt time.Time
analysisLastError string
analysisLastPersisted time.Time
analysisAggregateSeq uint64
analysisLearningScans analysisLearningScanAggregate
analysisEmbeddingBatches analysisEmbeddingBatchAggregate
processID string
analysisDetailMu sync.Mutex
analysisDetailVersion uint64
analysisDetailCache DetailedGraphAnalysis
researchTaskMu sync.Mutex
}
func ID(parts ...string) string {
h := sha256.Sum256([]byte(strings.Join(parts, "\x00")))
return hex.EncodeToString(h[:12])
}
func EdgeID(source, target, typ, origin string) string {
return ID("edge", source, target, typ, origin)
}
func (s *Store) markNodeDirtyLocked(id string) {
delete(s.deletedNodes, id)
s.dirtyNodes[id] = s.version
}
func (s *Store) markEdgeDirtyLocked(id string) {
delete(s.deletedEdges, id)
s.dirtyEdges[id] = s.version
}
func (s *Store) markVectorDirtyLocked(id string) {
delete(s.deletedVectors, id)
s.dirtyVectors[id] = s.version
}
func (s *Store) markNodeDeletedLocked(id string) {
delete(s.dirtyNodes, id)
s.deletedNodes[id] = s.version
}
func (s *Store) markEdgeDeletedLocked(id string) {
delete(s.dirtyEdges, id)
s.deletedEdges[id] = s.version
}
func (s *Store) markVectorDeletedLocked(id string) {
delete(s.dirtyVectors, id)
s.deletedVectors[id] = s.version
}
func (s *Store) Dirty() bool {
s.mu.RLock()
defer s.mu.RUnlock()
return len(s.dirtyNodes)+len(s.dirtyEdges)+len(s.dirtyVectors)+len(s.deletedNodes)+len(s.deletedEdges)+len(s.deletedVectors) > 0 || s.embeddingMetaGeneration > 0 || s.version != s.persistedVersion
}
// ConfigureEmbeddingModel records the model name before Ollama health data is
// available. A stored digest remains intact when the model name is unchanged.
func (s *Store) ConfigureEmbeddingModel(modelName string) int {
return s.ConfigureEmbeddingIdentity(modelName, "")
}
// ConfigureEmbeddingIdentity records the model and, once known, its Ollama
// digest. Importing a database built with another model identity automatically
// invalidates all vectors while retaining nodes and edges for selective
// relearning.
func (s *Store) ConfigureEmbeddingIdentity(modelName, digest string) int {
modelName = strings.TrimSpace(modelName)
digest = strings.TrimSpace(digest)
if modelName == "" {
return 0
}
s.mu.Lock()
defer s.mu.Unlock()
modelChanged := s.embeddingModel != "" && s.embeddingModel != modelName
digestChanged := digest != "" && s.embeddingDigest != "" && s.embeddingDigest != digest
metadataChanged := s.embeddingModel != modelName || (digest != "" && s.embeddingDigest != digest)
if !metadataChanged {
return 0
}
s.version++
removed := 0
if modelChanged || digestChanged {
for id, vector := range s.vectors {
label := ""
if node, ok := s.nodes[id]; ok {
label = node.Label
}
delete(s.vectors, id)
s.countVectorDeletedLocked()
s.recordChangeLocked(vectorChange(id, "deleted", len(vector), label))
s.deletedVectors[id] = s.version
delete(s.dirtyVectors, id)
removed++
}
}
s.embeddingModel = modelName
if modelChanged {
s.embeddingDigest = digest
} else if digest != "" {
s.embeddingDigest = digest
}
s.embeddingMetaGeneration = s.version
return removed
}
func (s *Store) UpsertNode(n model.Node) { _ = s.UpsertNodeWithStats(n) }
// UpsertNodeWithStats performs the same mutation as UpsertNode and returns only
// the mutation caused by this call. It is intentionally independent from the
// store-wide mutation counters so concurrent workflows can be attributed
// correctly in the analysis dashboard.
func (s *Store) UpsertNodeWithStats(n model.Node) MutationStats {
s.mu.Lock()
defer s.mu.Unlock()
var stats MutationStats
old, existed := s.nodes[n.ID]
if n.UpdatedAt.IsZero() {
n.UpdatedAt = time.Now().UTC()
}
if n.Weight == 0 {
n.Weight = 1
}
if n.X == 0 && n.Y == 0 && n.Z == 0 {
n.X, n.Y, n.Z = position(n.ID, n.Categories)
}
s.nodes[n.ID] = n
if existed {
s.countNodeUpdatedLocked()
stats.NodesUpdated++
} else {
s.countNodeCreatedLocked()
stats.NodesCreated++
}
s.version++
if existed {
s.recordChangeLocked(nodeUpdateChange(old, n))
} else {
s.recordChangeLocked(nodeChange(n, "created"))
}
s.markNodeDirtyLocked(n.ID)
return stats
}
func (s *Store) UpsertEdge(e model.Edge) { _ = s.UpsertEdgeWithStats(e) }
// UpsertEdgeWithStats is the causally attributable variant of UpsertEdge.
func (s *Store) UpsertEdgeWithStats(e model.Edge) MutationStats {
s.mu.Lock()
defer s.mu.Unlock()
var stats MutationStats
now := time.Now().UTC()
if e.ID == "" {
e.ID = EdgeID(e.Source, e.Target, e.Type, e.Origin)
}
old, existed := s.edges[e.ID]
if e.CreatedAt.IsZero() {
if existed {
e.CreatedAt = old.CreatedAt
} else {
e.CreatedAt = now
}
}
e.UpdatedAt = now
if e.Weight == 0 {
e.Weight = 1
}
s.edges[e.ID] = e
if existed {
s.countEdgeUpdatedLocked()
stats.EdgesUpdated++
} else {
s.countEdgeCreatedLocked()
stats.EdgesCreated++
}
s.version++
if existed {
s.recordChangeLocked(edgeUpdateChange(old, e))
} else {
s.recordChangeLocked(edgeChange(e, "created"))
}
s.markEdgeDirtyLocked(e.ID)
return stats
}
func (s *Store) HasEdgeBetween(a, b string) bool {
s.mu.RLock()
defer s.mu.RUnlock()
for _, e := range s.edges {
if (e.Source == a && e.Target == b) || (e.Source == b && e.Target == a) {
return true
}
}
return false
}
func (s *Store) GetNode(id string) (model.Node, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
n, ok := s.nodes[id]
return n, ok
}
func (s *Store) LookupExternal(id string) (model.Node, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
for _, n := range s.nodes {
if strings.EqualFold(n.ExternalID, id) {
return n, true
}
}
return model.Node{}, false
}
func (s *Store) SetVector(id string, v []float64) { _ = s.SetVectorWithStats(id, v) }
// SetVectorWithStats returns the exact vector mutation produced by this call.
func (s *Store) SetVectorWithStats(id string, v []float64) MutationStats {
s.mu.Lock()
defer s.mu.Unlock()
var stats MutationStats
converted := make([]float32, len(v))
for i, value := range v {
converted[i] = float32(value)
}
old, ok := s.vectors[id]
if ok && float32SlicesEqual(old, converted) {
return stats
}
s.vectors[id] = converted
if ok {
s.countVectorUpdatedLocked()
stats.VectorsUpdated++
} else {
s.countVectorCreatedLocked()
stats.VectorsCreated++
}
s.version++
label := ""
if node, exists := s.nodes[id]; exists {
label = node.Label
}
if ok {
s.recordChangeLocked(vectorRecalculatedChange(id, len(old), len(converted), label))
} else {
s.recordChangeLocked(vectorChange(id, "created", len(converted), label))
}
s.markVectorDirtyLocked(id)
return stats
}
func (s *Store) Vector(id string) ([]float64, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
v, ok := s.vectors[id]
if !ok {
return nil, false
}
out := make([]float64, len(v))
for i, value := range v {
out[i] = float64(value)
}
return out, true
}
func (s *Store) CountVectorsByDimension(dim int) int {
s.mu.RLock()
defer s.mu.RUnlock()
count := 0
for _, v := range s.vectors {
if len(v) == dim {
count++
}
}
return count
}
func (s *Store) ClearVectorsByDimension(dim int) int {
s.mu.Lock()
defer s.mu.Unlock()
removed := 0
for id, v := range s.vectors {
if len(v) == dim {
label := ""
if node, ok := s.nodes[id]; ok {
label = node.Label
}
delete(s.vectors, id)
s.countVectorDeletedLocked()
removed++
s.version++
s.recordChangeLocked(vectorChange(id, "deleted", len(v), label))
s.markVectorDeletedLocked(id)
}
}
return removed
}
func (s *Store) NodesForEmbedding() []model.Node {
return s.NodesForEmbeddingScoped(NodeFilter{})
}
func (s *Store) NodesForEmbeddingFiltered(sources []string) []model.Node {
return s.NodesForEmbeddingScoped(NodeFilter{Sources: sources})
}
func (s *Store) NodesForEmbeddingScoped(filter NodeFilter) []model.Node {
return s.nodesForEmbeddingKinds(filter, map[string]struct{}{"knowledge": {}, "ai-think": {}, "external": {}})
}
// KnowledgeNodesForEmbeddingScoped limits periodic KB learning to the nodes
// that are actually owned by the knowledge scanner. External research/security
// nodes are embedded by their own workflows and are only picked up by the
// global fallback-repair path when necessary.
func (s *Store) KnowledgeNodesForEmbeddingScoped(filter NodeFilter) []model.Node {
return s.nodesForEmbeddingKinds(filter, map[string]struct{}{"knowledge": {}, "ai-think": {}})
}
func (s *Store) nodesForEmbeddingKinds(filter NodeFilter, kinds map[string]struct{}) []model.Node {
s.mu.RLock()
defer s.mu.RUnlock()
out := []model.Node{}
for _, n := range s.nodes {
if _, ok := kinds[n.Kind]; !ok {
continue
}
if !filter.Matches(n) {
continue
}
if _, ok := s.vectors[n.ID]; !ok {
out = append(out, n)
}
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return out
}
func isRuntimeArticleProvenanceEdge(edge model.Edge) bool {
if edge.Origin != "knowledge-staging" && edge.Origin != "knowledge-synthesis" {
return false
}
return edge.Type == "synthesized_from" || edge.Type == "grounded_by" || strings.HasPrefix(edge.Type, "proposes_")
}
func (s *Store) KnowledgeNodes() []model.Node {
s.mu.RLock()
defer s.mu.RUnlock()
out := []model.Node{}
for _, n := range s.nodes {
if n.Kind == "knowledge" || n.Kind == "ai-think" {
out = append(out, n)
}
}
return out
}
func (s *Store) ReplaceOrigins(origins []string, nodes []model.Node, edges []model.Edge) {
_ = s.ReplaceOriginsWithStats(origins, nodes, edges)
}
// ReplaceOriginsWithStats reconciles managed origins and returns only the mutations
// caused by this reconciliation. This allows callers to report causal workflow
// changes even while other graph writers are active.
func (s *Store) ReplaceOriginsWithStats(origins []string, nodes []model.Node, edges []model.Edge) MutationStats {
var stats MutationStats
originSet := make(map[string]struct{}, len(origins))
for _, origin := range origins {
originSet[origin] = struct{}{}
}
now := time.Now().UTC()
incomingNodes := make(map[string]model.Node, len(nodes))
for _, node := range nodes {
if node.Weight == 0 {
node.Weight = 1
}
normalizeNodeCollections(&node)
incomingNodes[node.ID] = node
}
incomingEdges := make(map[string]model.Edge, len(edges))
for _, edge := range edges {
if edge.ID == "" {
edge.ID = EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin)
}
if edge.Weight == 0 {
edge.Weight = 1
}
normalizeEdgeCollections(&edge)
incomingEdges[edge.ID] = edge
}
s.mu.Lock()
defer s.mu.Unlock()
// Remove records from managed origins that disappeared from the source.
for id, old := range s.nodes {
if _, managed := originSet[old.Origin]; !managed {
continue
}
if _, present := incomingNodes[id]; present {
continue
}
delete(s.nodes, id)
s.countNodeDeletedLocked()
stats.NodesDeleted++
// A genuinely deleted node must not leave runtime provenance or any
// other unmanaged edge dangling. Re-import preservation applies only
// while the article node itself remains present.
for edgeID, edge := range s.edges {
if edge.Source != id && edge.Target != id {
continue
}
delete(s.edges, edgeID)
s.countEdgeDeletedLocked()
stats.EdgesDeleted++
s.version++
s.recordChangeLocked(edgeChange(edge, "deleted"))
s.markEdgeDeletedLocked(edgeID)
}
if _, hadVector := s.vectors[id]; hadVector {
vector := s.vectors[id]
delete(s.vectors, id)
s.countVectorDeletedLocked()
stats.VectorsDeleted++
s.version++
s.recordChangeLocked(vectorChange(id, "deleted", len(vector), old.Label))
s.markVectorDeletedLocked(id)
}
s.version++
s.recordChangeLocked(nodeChange(old, "deleted"))
s.markNodeDeletedLocked(id)
}
for id, old := range s.edges {
if _, managed := originSet[old.Origin]; !managed {
continue
}
// Runtime article provenance used knowledge-staging before the dedicated
// knowledge-synthesis origin existed. Those edges are not file-owned and
// must survive a staging directory reconciliation. Keeping them here also
// provides an in-place migration path for existing graphs.
if isRuntimeArticleProvenanceEdge(old) {
continue
}
if _, present := incomingEdges[id]; present {
continue
}
delete(s.edges, id)
s.countEdgeDeletedLocked()
stats.EdgesDeleted++
s.version++
s.recordChangeLocked(edgeChange(old, "deleted"))
s.markEdgeDeletedLocked(id)
}
// Reconcile nodes instead of deleting and recreating every row on each scan.
// This is what makes the scheduled SQLite flush truly incremental for an
// unchanged knowledge base.
for id, incoming := range incomingNodes {
old, existed := s.nodes[id]
if existed && incoming.Origin == "knowledge-staging" && old.Kind == "ai-think" {
if incoming.Metadata == nil {
incoming.Metadata = map[string]any{}
}
for _, key := range []string{"subtype", "action", "target_article_id", "target_node_id", "generation_depth", "confidence", "source_node_ids", "productive_source_count", "ai_source_count", "production_ratio", "source_fingerprint", "synthesis_model", "review_model", "pipeline"} {
if _, present := incoming.Metadata[key]; present {
continue
}
if value, present := old.Metadata[key]; present {
incoming.Metadata[key] = value
}
}
}
if incoming.UpdatedAt.IsZero() {
if existed && !old.UpdatedAt.IsZero() {
incoming.UpdatedAt = old.UpdatedAt
} else {
incoming.UpdatedAt = now
}
}
if incoming.X == 0 && incoming.Y == 0 && incoming.Z == 0 {
if existed && (old.X != 0 || old.Y != 0 || old.Z != 0) {
incoming.X, incoming.Y, incoming.Z = old.X, old.Y, old.Z
} else {
incoming.X, incoming.Y, incoming.Z = position(incoming.ID, incoming.Categories)
}
}
oldEmbeddingFingerprint := ""
if existed {
oldEmbeddingFingerprint = embeddingFingerprint(old)
}
if existed && nodesEquivalent(old, incoming) {
continue
}
s.nodes[id] = incoming
if existed {
s.countNodeUpdatedLocked()
stats.NodesUpdated++
} else {
s.countNodeCreatedLocked()
stats.NodesCreated++
}
s.version++
if existed {
s.recordChangeLocked(nodeUpdateChange(old, incoming))
} else {
s.recordChangeLocked(nodeChange(incoming, "created"))
}
s.markNodeDirtyLocked(id)
// Embeddings depend on text/categories/keywords, not on display
// coordinates or unrelated metadata. Only invalidate a vector when its
// actual embedding input changed.
if existed && oldEmbeddingFingerprint != embeddingFingerprint(incoming) {
if vector, hadVector := s.vectors[id]; hadVector {
delete(s.vectors, id)
s.countVectorDeletedLocked()
stats.VectorsDeleted++
s.version++
s.recordChangeLocked(vectorChange(id, "deleted", len(vector), incoming.Label))
s.markVectorDeletedLocked(id)
}
}
}
// Reconcile deterministic source edges. Timestamps are retained for an
// unchanged edge so periodic scans do not produce needless writes.
for id, incoming := range incomingEdges {
old, existed := s.edges[id]
if existed && edgesEquivalentIgnoringTimestamps(old, incoming) {
continue
}
if incoming.CreatedAt.IsZero() {
if existed && !old.CreatedAt.IsZero() {
incoming.CreatedAt = old.CreatedAt
} else {
incoming.CreatedAt = now
}
}
incoming.UpdatedAt = now
s.edges[id] = incoming
if existed {
s.countEdgeUpdatedLocked()
stats.EdgesUpdated++
} else {
s.countEdgeCreatedLocked()
stats.EdgesCreated++
}
s.version++
if existed {
s.recordChangeLocked(edgeUpdateChange(old, incoming))
} else {
s.recordChangeLocked(edgeChange(incoming, "created"))
}
s.markEdgeDirtyLocked(id)
}
// Remove any remaining edge whose endpoint no longer exists. This includes
// AI-derived edges that referred to a source note removed from the KB.
for id, edge := range s.edges {
if _, ok := s.nodes[edge.Source]; !ok {
delete(s.edges, id)
s.countEdgeDeletedLocked()
stats.EdgesDeleted++
s.version++
s.recordChangeLocked(edgeChange(edge, "deleted"))
s.markEdgeDeletedLocked(id)
continue
}
if _, ok := s.nodes[edge.Target]; !ok {
delete(s.edges, id)
s.countEdgeDeletedLocked()
stats.EdgesDeleted++
s.version++
s.recordChangeLocked(edgeChange(edge, "deleted"))
s.markEdgeDeletedLocked(id)
}
}
return stats
}
func embeddingFingerprint(node model.Node) string {
return node.Label + "\x00" + node.Summary + "\x00" + strings.Join(node.Categories, "\x00") + "\x00" + strings.Join(node.Keywords, "\x00")
}
func normalizeNodeCollections(node *model.Node) {
if node.Categories == nil {
node.Categories = []string{}
}
if node.Keywords == nil {
node.Keywords = []string{}
}
if node.Metadata == nil {
node.Metadata = map[string]any{}
}
}
func normalizeEdgeCollections(edge *model.Edge) {
if edge.Evidence == nil {
edge.Evidence = []model.Evidence{}
}
if edge.Metadata == nil {
edge.Metadata = map[string]any{}
}
}
func nodesEquivalent(a, b model.Node) bool {
if a.ID != b.ID || a.Kind != b.Kind || a.Label != b.Label ||
a.Summary != b.Summary || a.Status != b.Status || a.Origin != b.Origin ||
a.ExternalID != b.ExternalID || a.URI != b.URI || a.Weight != b.Weight ||
a.X != b.X || a.Y != b.Y || a.Z != b.Z {
return false
}
return JSONEquivalent(a.Categories, b.Categories) &&
JSONEquivalent(a.Keywords, b.Keywords) &&
JSONEquivalent(a.Metadata, b.Metadata)
}
func edgesEquivalentIgnoringTimestamps(a, b model.Edge) bool {
if a.ID != b.ID || a.Source != b.Source || a.Target != b.Target ||
a.Type != b.Type || a.Origin != b.Origin || a.Status != b.Status ||
a.Confidence != b.Confidence || a.Weight != b.Weight ||
a.Explanation != b.Explanation {
return false
}
return JSONEquivalent(a.Evidence, b.Evidence) && JSONEquivalent(a.Metadata, b.Metadata)
}
func JSONEquivalent(a, b any) bool {
left, leftErr := json.Marshal(a)
right, rightErr := json.Marshal(b)
return leftErr == nil && rightErr == nil && string(left) == string(right)
}
func (s *Store) Version() uint64 {
s.mu.RLock()
defer s.mu.RUnlock()
return s.version
}
func (s *Store) Counts() (nodes, edges int, version uint64) {
s.mu.RLock()
defer s.mu.RUnlock()
nodes = len(s.nodes)
for _, edge := range s.edges {
if edge.Status != "rejected" {
edges++
}
}
return nodes, edges, s.version
}
func (s *Store) IdleNode(seed int64) (model.Node, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
if len(s.nodes) == 0 {
return model.Node{}, false
}
index := int(seed % int64(len(s.nodes)))
if index < 0 {
index = -index
}
for _, node := range s.nodes {
if index == 0 {
return node, true
}
index--
}
return model.Node{}, false
}
func (s *Store) Snapshot() model.Snapshot {
s.mu.RLock()
defer s.mu.RUnlock()
n := make([]model.Node, 0, len(s.nodes))
e := make([]model.Edge, 0, len(s.edges))
for _, x := range s.nodes {
n = append(n, x)
}
for _, x := range s.edges {
if x.Status == "rejected" {
continue
}
e = append(e, x)
}
sort.Slice(n, func(i, j int) bool { return n[i].ID < n[j].ID })
sort.Slice(e, func(i, j int) bool { return e[i].ID < e[j].ID })
return model.Snapshot{Version: s.version, Nodes: n, Edges: e, UpdatedAt: time.Now().UTC()}
}
func (s *Store) Similar(query []float64, limit int) []model.Hit {
return s.SimilarFiltered(query, limit, NodeFilter{})
}
func (s *Store) SimilarFiltered(query []float64, limit int, filter NodeFilter) []model.Hit {
s.mu.RLock()
defer s.mu.RUnlock()
hits := []model.Hit{}
for id, v := range s.vectors {
n, ok := s.nodes[id]
if !ok || (n.Kind != "knowledge" && n.Kind != "ai-think" && n.Kind != "external") || !filter.Matches(n) {
continue
}
score := cosineMixed(query, v)
hits = append(hits, model.Hit{NodeID: id, Label: n.Label, Score: score, Kind: n.Kind, Status: n.Status})
}
sort.Slice(hits, func(i, j int) bool { return hits[i].Score > hits[j].Score })
if limit > 0 && len(hits) > limit {
hits = hits[:limit]
}
return hits
}
// NextPair searches a bounded rotating window of anchor nodes instead of
// comparing the complete graph on every AI-THINK cycle. This keeps candidate
// selection responsive even for tens of thousands of knowledge nodes while the
// rotating cursor eventually visits the complete corpus.
func (s *Store) NextPair(min float64, anchorLimit int) (model.Node, model.Node, float64, bool, int) {
return s.NextPairFiltered(min, anchorLimit, nil)
}
func (s *Store) NextPairFiltered(min float64, anchorLimit int, sources []string) (model.Node, model.Node, float64, bool, int) {
return s.NextPairFilteredDepth(min, anchorLimit, sources, 0)
}
func (s *Store) NextPairFilteredDepth(min float64, anchorLimit int, sources []string, maxAIDepth int) (model.Node, model.Node, float64, bool, int) {
return s.NextPairScopedDepth(min, anchorLimit, NodeFilter{Sources: sources}, maxAIDepth)
}
func (s *Store) NextPairScopedDepth(min float64, anchorLimit int, filter NodeFilter, maxAIDepth int) (model.Node, model.Node, float64, bool, int) {
s.mu.Lock()
defer s.mu.Unlock()
nodes := make([]model.Node, 0, len(s.nodes))
for _, n := range s.nodes {
if n.Kind != "knowledge" && n.Kind != "ai-think" {
continue
}
if !filter.Matches(n) {
continue
}
if n.Kind == "ai-think" && maxAIDepth > 0 && graphNodeGenerationDepth(n) >= maxAIDepth {
continue
}
if v, ok := s.vectors[n.ID]; ok && len(v) > 0 {
nodes = append(nodes, n)
}
}
if len(nodes) < 2 {
return model.Node{}, model.Node{}, 0, false, 0
}
sort.Slice(nodes, func(i, j int) bool { return nodes[i].ID < nodes[j].ID })
if anchorLimit <= 0 || anchorLimit > len(nodes) {
anchorLimit = len(nodes)
}
blocked := make(map[string]struct{}, len(s.edges))
for _, e := range s.edges {
blocked[pairKey(e.Source, e.Target)] = struct{}{}
}
start := s.pairCursor % len(nodes)
best := -1.0
var a, b model.Node
comparisons := 0
for step := 0; step < anchorLimit; step++ {
i := (start + step) % len(nodes)
left := nodes[i]
lv := s.vectors[left.ID]
for j := 0; j < len(nodes); j++ {
if i == j {
continue
}
right := nodes[j]
if left.Kind == "ai-think" && right.Kind == "ai-think" {
continue
}
if _, exists := blocked[pairKey(left.ID, right.ID)]; exists {
continue
}
rv := s.vectors[right.ID]
if len(lv) != len(rv) {
continue
}
comparisons++
score := cosine32(lv, rv)
if score >= min && score > best {
best = score
a, b = left, right
}
}
}
s.pairCursor = (start + anchorLimit) % len(nodes)
return a, b, best, best >= 0, comparisons
}
func pairKey(a, b string) string {
if a > b {
a, b = b, a
}
return a + "\x00" + b
}
func (s *Store) BestPair(min float64) (model.Node, model.Node, float64, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
nodes := []model.Node{}
for _, n := range s.nodes {
if n.Kind == "knowledge" || n.Kind == "ai-think" {
if _, ok := s.vectors[n.ID]; ok {
nodes = append(nodes, n)
}
}
}
best := -1.0
var a, b model.Node
for i := 0; i < len(nodes); i++ {
for j := i + 1; j < len(nodes); j++ {
if edgeBetweenLocked(s.edges, nodes[i].ID, nodes[j].ID) {
continue
}
score := cosine32(s.vectors[nodes[i].ID], s.vectors[nodes[j].ID])
if score >= min && score > best {
best = score
a = nodes[i]
b = nodes[j]
}
}
}
return a, b, best, best >= 0
}
func (s *Store) ConnectingEdges(ids []string) []string {
set := map[string]bool{}
for _, id := range ids {
set[id] = true
}
s.mu.RLock()
defer s.mu.RUnlock()
var out []string
for id, e := range s.edges {
if set[e.Source] && set[e.Target] {
out = append(out, id)
}
}
return out
}
func graphNodeGenerationDepth(n model.Node) int {
if n.Kind != "ai-think" {
return 0
}
value, ok := n.Metadata["generation_depth"]
if !ok {
return 1
}
switch typed := value.(type) {
case int:
return typed
case int64:
return int(typed)
case float64:
return int(typed)
case json.Number:
value, _ := typed.Int64()
return int(value)
default:
return 1
}
}
func edgeBetweenLocked(edges map[string]model.Edge, a, b string) bool {
for _, e := range edges {
if (e.Source == a && e.Target == b) || (e.Source == b && e.Target == a) {
return true
}
}
return false
}
func float32SlicesEqual(a, b []float32) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
func cosine32(a, b []float32) float64 {
if len(a) == 0 || len(a) != len(b) {
return 0
}
var dot, aa, bb float64
for i := range a {
av, bv := float64(a[i]), float64(b[i])
dot += av * bv
aa += av * av
bb += bv * bv
}
if aa == 0 || bb == 0 {
return 0
}
return dot / (math.Sqrt(aa) * math.Sqrt(bb))
}
func cosineMixed(a []float64, b []float32) float64 {
if len(a) == 0 || len(a) != len(b) {
return 0
}
var dot, aa, bb float64
for i := range a {
bv := float64(b[i])
dot += a[i] * bv
aa += a[i] * a[i]
bb += bv * bv
}
if aa == 0 || bb == 0 {
return 0
}
return dot / (math.Sqrt(aa) * math.Sqrt(bb))
}
func position(id string, cats []string) (float64, float64, float64) {
seed := sha256.Sum256([]byte(id + "\x00" + strings.Join(cats, "|")))
u := func(i int) float64 { return float64(int(seed[i%len(seed)])) / 255 }
side := -1.0
if seed[0]%2 == 0 {
side = 1
}
biasY, biasZ := 0.0, 0.0
if len(cats) > 0 {
h := sha256.Sum256([]byte(cats[0]))
biasY = (float64(h[0])/255 - .5) * .9
biasZ = (float64(h[1])/255 - .5) * .65
}
for i := 0; i < 16; i++ {
x := side * (0.08 + u(1+i)*0.72)
y := biasY*.32 + (u(2+i)-.5)*1.18
z := biasZ*.28 + (u(3+i)-.5)*.94
if insideBrainShape(x, y, z) {
return x, y, z
}
}
return side * .34, biasY * .22, biasZ * .2
}
func insideBrainShape(x, y, z float64) bool {
if math.Abs(x) < .045 && y > -.58 && y < .42 {
return false
}
if y < -.76 || y > .82 {
return false
}
taperY := y + math.Abs(z)*.10 - math.Max(0, math.Abs(x)-.58)*.18
lx := (x + .35) / .58
rx := (x - .35) / .58
ny := taperY / .76
nz := z / .58
left := lx*lx+ny*ny+nz*nz <= 1
right := rx*rx+ny*ny+nz*nz <= 1
return left || right
}
func (s *Store) Analyze() model.GraphAnalysis {
s.mu.RLock()
defer s.mu.RUnlock()
analysis := model.GraphAnalysis{NodeCount: len(s.nodes)}
degree := make(map[string]int, len(s.nodes))
knowledgeLinked := make(map[string]bool)
parent := make(map[string]string, len(s.nodes))
for id, n := range s.nodes {
parent[id] = id
if n.Status == "staging" {
analysis.StagingNodes++
}
if n.Kind == "ai-think" {
analysis.AIThinkNodes++
}
if n.Kind == "external" {
analysis.ExternalNodes++
}
}
var find func(string) string
find = func(x string) string {
p := parent[x]
if p != x {
parent[x] = find(p)
}
return parent[x]
}
union := func(a, b string) {
ra, rb := find(a), find(b)
if ra != rb {
parent[rb] = ra
}
}
for _, e := range s.edges {
if e.Status == "rejected" {
continue
}
if _, ok := s.nodes[e.Source]; !ok {
continue
}
if _, ok := s.nodes[e.Target]; !ok {
continue
}
analysis.EdgeCount++
degree[e.Source]++
degree[e.Target]++
union(e.Source, e.Target)
if e.Origin == "ai-inference" {
analysis.AIEdges++
}
if e.Type == "contradicts" {
analysis.Contradictions++
}
a, b := s.nodes[e.Source], s.nodes[e.Target]
aKnowledge := a.Kind == "knowledge" || a.Kind == "ai-think"
bKnowledge := b.Kind == "knowledge" || b.Kind == "ai-think"
// Knowledge connectivity is undirected for readiness purposes. Research and
// security evidence often point external -> knowledge, while article links
// point knowledge -> knowledge. Both directions must count consistently.
if aKnowledge && (bKnowledge || b.Kind == "external") {
knowledgeLinked[a.ID] = true
}
if bKnowledge && (aKnowledge || a.Kind == "external") {
knowledgeLinked[b.ID] = true
}
}
roots := map[string]bool{}
for id, n := range s.nodes {
roots[find(id)] = true
if (n.Kind == "knowledge" || n.Kind == "ai-think") && !knowledgeLinked[id] {
analysis.KnowledgeOrphans++
}
}
analysis.Components = len(roots)
hubs := make([]model.Hub, 0, len(degree))
for id, d := range degree {
n := s.nodes[id]
hubs = append(hubs, model.Hub{NodeID: id, Label: n.Label, Kind: n.Kind, Degree: d})
}
sort.Slice(hubs, func(i, j int) bool {
if hubs[i].Degree == hubs[j].Degree {
return hubs[i].Label < hubs[j].Label
}
return hubs[i].Degree > hubs[j].Degree
})
if len(hubs) > 8 {
hubs = hubs[:8]
}
analysis.TopHubs = hubs
return analysis
}