All checks were successful
release-tag / release-image (push) Successful in 2m32s
1535 lines
74 KiB
Go
1535 lines
74 KiB
Go
package engine
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"math"
|
|
"path/filepath"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
"unicode"
|
|
|
|
"github.com/local/glpi-neural-brain/internal/activity"
|
|
"github.com/local/glpi-neural-brain/internal/config"
|
|
"github.com/local/glpi-neural-brain/internal/glpi"
|
|
"github.com/local/glpi-neural-brain/internal/graph"
|
|
"github.com/local/glpi-neural-brain/internal/ingest"
|
|
"github.com/local/glpi-neural-brain/internal/model"
|
|
"github.com/local/glpi-neural-brain/internal/ollama"
|
|
"github.com/local/glpi-neural-brain/internal/persist"
|
|
"github.com/local/glpi-neural-brain/internal/research"
|
|
"github.com/local/glpi-neural-brain/internal/sourceagent"
|
|
"github.com/local/glpi-neural-brain/internal/workqueue"
|
|
)
|
|
|
|
var (
|
|
ErrNoCandidate = errors.New("no enrichment candidate")
|
|
ErrLearningDisabled = errors.New("learning is disabled")
|
|
ErrThinkingDisabled = errors.New("thinking is disabled")
|
|
)
|
|
|
|
type pendingArticleCandidate struct {
|
|
Seeds []model.Node
|
|
Relation model.RelationDecision
|
|
Research []model.ResearchResult
|
|
}
|
|
|
|
type EnrichOutcome struct {
|
|
Result string
|
|
Candidate bool
|
|
Created bool
|
|
RelationCreated bool
|
|
ArticleCreated bool
|
|
ArticleSkipped bool
|
|
Rejected bool
|
|
Comparisons int
|
|
CoarseComparisons int
|
|
IndexedNodes int
|
|
CandidatePool int
|
|
PendingArticle *pendingArticleCandidate
|
|
Mutations graph.MutationStats
|
|
}
|
|
|
|
type Engine struct {
|
|
Cfg config.Config
|
|
Graph *graph.Store
|
|
Broker *activity.Broker
|
|
Ollama *ollama.Client
|
|
Research *research.Client
|
|
Scanner *ingest.KnowledgeScanner
|
|
GLPIKB *ingest.GLPIKBSyncer
|
|
Persistence *persist.Coordinator
|
|
SourceInbox *sourceagent.Store
|
|
|
|
bootstrapReady chan struct{}
|
|
bootstrapOnce sync.Once
|
|
|
|
mu sync.Mutex
|
|
stateMu sync.RWMutex
|
|
lastScan time.Time
|
|
lastVectorGraph time.Time
|
|
lastVectorLayout time.Time
|
|
vectorMaintenanceStartedAt time.Time
|
|
bootstrapComplete bool
|
|
bootstrapCompletedAt time.Time
|
|
bootstrapError string
|
|
lastEnrich time.Time
|
|
lastAttempt time.Time
|
|
nextEnrich time.Time
|
|
ollamaOK bool
|
|
enrichRunning bool
|
|
enrichTrigger string
|
|
enrichResult string
|
|
enrichError string
|
|
enrichCycles uint64
|
|
enrichCreated uint64
|
|
enrichRejected uint64
|
|
relationsCreated uint64
|
|
articlesCreated uint64
|
|
articlesSkipped uint64
|
|
enrichRequests chan string
|
|
runtimeMu sync.RWMutex
|
|
runtime RuntimeSettings
|
|
runtimePath string
|
|
researchEvidenceMu sync.RWMutex
|
|
researchEvidenceCache map[string]researchEvidenceRecord
|
|
sharedWork *workqueue.Limiter
|
|
researchDedupeMu sync.Mutex
|
|
researchDedupe map[string]*researchDedupeEntry
|
|
researchDedupeGuardFiltered uint64
|
|
researchDedupeGuardPrimaryMismatch uint64
|
|
researchDedupeGuardFocusMismatch uint64
|
|
researchDedupeGuardEntityMismatch uint64
|
|
interactiveInflight atomic.Int64
|
|
autonomousWake chan struct{}
|
|
autonomousScanRequests chan string
|
|
sourceInboxWake chan struct{}
|
|
autonomousRunning bool
|
|
autonomousTaskID string
|
|
autonomousTaskTopic string
|
|
autonomousLastStarted time.Time
|
|
autonomousLastCompleted time.Time
|
|
autonomousLastError string
|
|
autonomousLastScanStarted time.Time
|
|
autonomousLastScanCompleted time.Time
|
|
autonomousLastScanTrigger string
|
|
autonomousLastScanCandidates int
|
|
autonomousLastScanCreated int
|
|
autonomousLastScanDecisions []autonomousOpportunityDecision
|
|
autonomousCompleted uint64
|
|
autonomousFailed uint64
|
|
autonomousEvidence uint64
|
|
autonomousArticles uint64
|
|
}
|
|
|
|
func New(cfg config.Config, g *graph.Store, b *activity.Broker) *Engine {
|
|
if b != nil && g != nil {
|
|
b.SetSink(g.RecordActivity)
|
|
}
|
|
if strings.TrimSpace(cfg.GLPIKBSource) == "" {
|
|
cfg.GLPIKBSource = "GLPI Knowledge Base"
|
|
}
|
|
if !cfg.RuntimeDefaultsConfigured {
|
|
cfg.LearningEnabled = true
|
|
cfg.ThinkingEnabled = true
|
|
cfg.DefaultView = "neural"
|
|
}
|
|
if cfg.EnrichBatchSize < 1 {
|
|
cfg.EnrichBatchSize = 1
|
|
}
|
|
if cfg.EnrichAnchors < 1 {
|
|
cfg.EnrichAnchors = 48
|
|
}
|
|
if cfg.ArticleMinSources < 2 {
|
|
cfg.ArticleMinSources = 3
|
|
}
|
|
if cfg.ArticleMaxSources < cfg.ArticleMinSources {
|
|
cfg.ArticleMaxSources = 8
|
|
}
|
|
if cfg.ArticleMinProductionRatio == 0 {
|
|
cfg.ArticleMinProductionRatio = .70
|
|
}
|
|
if cfg.ArticleMaxGenerationDepth < 1 {
|
|
cfg.ArticleMaxGenerationDepth = 2
|
|
}
|
|
if cfg.ArticleMinConfidence == 0 {
|
|
cfg.ArticleMinConfidence = .74
|
|
}
|
|
if cfg.ArticleMinTextChars < 1 {
|
|
cfg.ArticleMinTextChars = 180
|
|
}
|
|
if cfg.ArticleMinAnswerChars < 1 {
|
|
cfg.ArticleMinAnswerChars = 420
|
|
}
|
|
if cfg.ArticleMaxResearchQueries < 1 {
|
|
cfg.ArticleMaxResearchQueries = 6
|
|
}
|
|
if cfg.ArticleResearchResults < 1 {
|
|
cfg.ArticleResearchResults = 12
|
|
}
|
|
if cfg.ArticleResearchRounds < 1 {
|
|
cfg.ArticleResearchRounds = 3
|
|
}
|
|
if cfg.ArticleResearchFetchResults < 1 {
|
|
cfg.ArticleResearchFetchResults = 6
|
|
}
|
|
if cfg.ArticleResearchFetchResults > cfg.ArticleResearchResults {
|
|
cfg.ArticleResearchFetchResults = cfg.ArticleResearchResults
|
|
}
|
|
if cfg.ArticleAdaptiveInitialQueries < 1 {
|
|
cfg.ArticleAdaptiveInitialQueries = 2
|
|
}
|
|
if cfg.ArticleAdaptiveInitialFetch < 1 {
|
|
cfg.ArticleAdaptiveInitialFetch = 3
|
|
}
|
|
if cfg.ArticleAdaptiveInitialFetch > cfg.ArticleResearchFetchResults {
|
|
cfg.ArticleAdaptiveInitialFetch = cfg.ArticleResearchFetchResults
|
|
}
|
|
if cfg.ArticleResearchExplorationResults > cfg.ArticleResearchFetchResults {
|
|
cfg.ArticleResearchExplorationResults = cfg.ArticleResearchFetchResults
|
|
}
|
|
if cfg.ArticleResearchPrefetchMinRelevance <= 0 {
|
|
cfg.ArticleResearchPrefetchMinRelevance = .25
|
|
}
|
|
if cfg.ArticleResearchMinRelevance <= 0 {
|
|
cfg.ArticleResearchMinRelevance = .55
|
|
}
|
|
if cfg.ArticleResearchPrefetchMinRelevance > cfg.ArticleResearchMinRelevance {
|
|
cfg.ArticleResearchPrefetchMinRelevance = cfg.ArticleResearchMinRelevance
|
|
}
|
|
if cfg.ArticleResearchMinQuality <= 0 {
|
|
cfg.ArticleResearchMinQuality = .35
|
|
}
|
|
if cfg.ArticleResearchPageMaxBytes < 1 {
|
|
cfg.ArticleResearchPageMaxBytes = 2 << 20
|
|
}
|
|
if cfg.ArticleResearchPageMaxChars < 1 {
|
|
cfg.ArticleResearchPageMaxChars = 14000
|
|
}
|
|
if cfg.ArticleResearchFetchTimeout < time.Second {
|
|
cfg.ArticleResearchFetchTimeout = 20 * time.Second
|
|
}
|
|
if strings.TrimSpace(cfg.ArticleLanguage) == "" {
|
|
cfg.ArticleLanguage = "de-DE"
|
|
}
|
|
if strings.TrimSpace(cfg.ArticleSynthesisModel) == "" {
|
|
cfg.ArticleSynthesisModel = cfg.ChatModel
|
|
}
|
|
if strings.TrimSpace(cfg.ArticleReviewModel) == "" {
|
|
cfg.ArticleReviewModel = cfg.ChatModel
|
|
}
|
|
if cfg.ResearchDedupeThreshold <= 0 {
|
|
cfg.ResearchDedupeThreshold = .92
|
|
}
|
|
if cfg.ResearchDedupeTTL <= 0 {
|
|
cfg.ResearchDedupeTTL = 45 * time.Minute
|
|
}
|
|
if cfg.ResearchOllamaMaxInflight < 1 {
|
|
cfg.ResearchOllamaMaxInflight = 2
|
|
}
|
|
if cfg.ResearchOllamaQueueSize < 1 {
|
|
cfg.ResearchOllamaQueueSize = 64
|
|
}
|
|
if cfg.AutonomousResearchInterval < time.Minute {
|
|
cfg.AutonomousResearchInterval = 30 * time.Minute
|
|
}
|
|
if cfg.AutonomousResearchTasksPerCycle < 1 {
|
|
cfg.AutonomousResearchTasksPerCycle = 1
|
|
}
|
|
if cfg.AutonomousResearchMaxTasksPerDay < 1 {
|
|
cfg.AutonomousResearchMaxTasksPerDay = 12
|
|
}
|
|
if cfg.AutonomousResearchMaxQueriesPerTask < 1 {
|
|
cfg.AutonomousResearchMaxQueriesPerTask = 6
|
|
}
|
|
if cfg.AutonomousResearchMaxPagesPerTask < 1 {
|
|
cfg.AutonomousResearchMaxPagesPerTask = 8
|
|
}
|
|
if cfg.AutonomousResearchMaxRounds < 1 {
|
|
cfg.AutonomousResearchMaxRounds = 3
|
|
}
|
|
if cfg.AutonomousResearchMinPriority <= 0 {
|
|
cfg.AutonomousResearchMinPriority = .65
|
|
}
|
|
if cfg.AutonomousResearchCooldown < time.Hour {
|
|
cfg.AutonomousResearchCooldown = 168 * time.Hour
|
|
}
|
|
if cfg.AutonomousResearchLease < 5*time.Minute {
|
|
cfg.AutonomousResearchLease = 45 * time.Minute
|
|
}
|
|
if cfg.AutonomousResearchMaxAttempts < 1 {
|
|
cfg.AutonomousResearchMaxAttempts = 3
|
|
}
|
|
if cfg.AutonomousResearchOpportunityLimit < 1 {
|
|
cfg.AutonomousResearchOpportunityLimit = 8
|
|
}
|
|
if cfg.SourceInboxSecurityBatchSize < 1 {
|
|
cfg.SourceInboxSecurityBatchSize = 2
|
|
}
|
|
if cfg.SourceInboxSecurityMinPriority <= 0 {
|
|
cfg.SourceInboxSecurityMinPriority = .58
|
|
}
|
|
if cfg.SourceInboxSecurityMinConfidence <= 0 {
|
|
cfg.SourceInboxSecurityMinConfidence = .72
|
|
}
|
|
if cfg.SourceInboxSecurityFetchMinChars < 160 {
|
|
cfg.SourceInboxSecurityFetchMinChars = 1800
|
|
}
|
|
if cfg.SourceInboxSecurityResearchResults < 1 {
|
|
cfg.SourceInboxSecurityResearchResults = 3
|
|
}
|
|
ollamaURLs := append([]string(nil), cfg.OllamaURLs...)
|
|
if len(ollamaURLs) == 0 && strings.TrimSpace(cfg.OllamaURL) != "" {
|
|
ollamaURLs = []string{cfg.OllamaURL}
|
|
}
|
|
nodes := make([]ollama.NodeConfig, 0, len(ollamaURLs))
|
|
for i, rawURL := range ollamaURLs {
|
|
name := fmt.Sprintf("ollama-%d", i+1)
|
|
if i < len(cfg.OllamaNodeNames) && strings.TrimSpace(cfg.OllamaNodeNames[i]) != "" {
|
|
name = strings.TrimSpace(cfg.OllamaNodeNames[i])
|
|
}
|
|
weight := 1
|
|
if i < len(cfg.OllamaNodeWeights) && cfg.OllamaNodeWeights[i] > 0 {
|
|
weight = cfg.OllamaNodeWeights[i]
|
|
}
|
|
nodes = append(nodes, ollama.NodeConfig{Name: name, URL: rawURL, Weight: weight})
|
|
}
|
|
clearedVectors := g.ConfigureEmbeddingModel(cfg.EmbeddingModel)
|
|
if clearedVectors > 0 && b != nil {
|
|
b.Publish(model.Activity{Type: "embedding.model_changed", Source: "brain", Phase: "learning", Message: fmt.Sprintf("Embedding-Modell geändert · %d Vektoren werden neu gelernt", clearedVectors), Strength: .65, Metadata: map[string]any{"embedding_model": cfg.EmbeddingModel, "cleared_vectors": clearedVectors}})
|
|
}
|
|
pool := ollama.NewPool(ollama.PoolConfig{
|
|
Nodes: nodes, RoutingMode: cfg.OllamaRoutingMode, NodeMaxInflight: cfg.OllamaNodeMaxInflight,
|
|
HealthInterval: cfg.OllamaHealthInterval, FailureCooldown: cfg.OllamaFailureCooldown,
|
|
RequestTimeout: cfg.OllamaRequestTimeout, FailoverEnabled: cfg.OllamaFailoverEnabled,
|
|
FailoverAttempts: cfg.OllamaFailoverAttempts, RequireSameModelDigest: cfg.OllamaRequireSameDigest,
|
|
RequireEmbeddingModel: cfg.OllamaRequireEmbeddingModel,
|
|
}, cfg.ChatModel, cfg.EmbeddingModel)
|
|
sharedWork := workqueue.New(cfg.ResearchOllamaMaxInflight, cfg.ResearchOllamaQueueSize)
|
|
pool.SetSharedLimiter(sharedWork)
|
|
persistence := persist.New(g, b, cfg.PersistInterval)
|
|
e := &Engine{Cfg: cfg, Graph: g, Broker: b, Ollama: pool, Persistence: persistence, vectorMaintenanceStartedAt: time.Now().UTC(), Scanner: &ingest.KnowledgeScanner{Graph: g, ProductionDirs: cfg.KnowledgeDirs, StagingDirs: cfg.StagingDirs, FullVerifyInterval: cfg.KnowledgeFullVerifyInterval}, bootstrapReady: make(chan struct{}), enrichRequests: make(chan string, 1), autonomousWake: make(chan struct{}, 1), autonomousScanRequests: make(chan string, 1), sourceInboxWake: make(chan struct{}, 1), runtimePath: filepath.Join(cfg.DataDir, "runtime-settings.json"), researchEvidenceCache: map[string]researchEvidenceRecord{}, sharedWork: sharedWork, researchDedupe: map[string]*researchDedupeEntry{}}
|
|
e.loadRuntimeSettings()
|
|
e.applyRuntimePerformance(e.RuntimeSettings())
|
|
if cfg.SearXNGURL != "" {
|
|
e.Research = research.New(cfg.SearXNGURL)
|
|
}
|
|
if cfg.GLPIKBEnabled {
|
|
client := glpi.New(cfg.GLPIURL, cfg.GLPIAPIVersion, cfg.GLPIClientID, cfg.GLPIClientSecret, cfg.GLPIUsername, cfg.GLPIPassword, cfg.GLPITimeout)
|
|
e.GLPIKB = ingest.NewGLPIKBSyncer(ingest.GLPIKBConfig{Enabled: true, Path: cfg.GLPIKBPath, Filter: cfg.GLPIKBFilter, Limit: cfg.GLPIKBLimit, SyncInterval: cfg.GLPIKBSyncInterval, Source: cfg.GLPIKBSource, CachePath: filepath.Join(cfg.DataDir, "glpi-kb-cache.json"), ShouldSync: e.LearningEnabled}, client, g, b, persistence)
|
|
}
|
|
if b != nil {
|
|
b.Publish(model.Activity{Type: "system.started", Source: "brain", Phase: "startup", Message: "Neural Brain wurde gestartet; das Analyseprotokoll zeichnet Läufe und Graphänderungen auf", Strength: .3, Metadata: map[string]any{"chat_model": cfg.ChatModel, "embedding_model": cfg.EmbeddingModel, "article_language": cfg.ArticleLanguage, "article_synthesis_model": cfg.ArticleSynthesisModel, "article_review_model": cfg.ArticleReviewModel, "article_review_repair_rounds": cfg.ArticleReviewRepairRounds, "article_pipeline": "adaptive_generate_review/v5-quality-gate-v12", "article_research_strategy": cfg.ArticleResearchStrategy, "cluster_article_batching": cfg.ClusterArticleBatching, "vector_graph_enabled": cfg.VectorGraphEnabled, "vector_graph_layout": cfg.VectorGraphLayout, "vector_graph_layout_mode": e.effectiveVectorLayoutMode(), "vector_graph_reevaluate_interval": cfg.VectorGraphReevaluateInterval.String(), "vector_graph_relax_layout": cfg.VectorGraphRelaxLayout, "vector_graph_layout_relax_effective": cfg.VectorGraphRelaxLayout && !cfg.VectorGraphLayout, "vector_graph_layout_relax_interval": cfg.VectorGraphLayoutRelaxInterval.String(), "article_cpu_quality_enabled": cfg.ArticleCPUQualityEnabled, "article_cpu_quality_agent_offload": cfg.ArticleCPUQualityAgentOffload, "research_ollama_max_inflight": cfg.ResearchOllamaMaxInflight, "research_ollama_queue_size": cfg.ResearchOllamaQueueSize, "speed_mode": e.RuntimeSettings().SpeedMode, "speed_cpu_tasks": e.RuntimeSettings().SpeedCPUWorkers, "speed_gpu_tasks": e.RuntimeSettings().SpeedGPUInflight, "graph_version": g.Version()}})
|
|
}
|
|
return e
|
|
}
|
|
func (e *Engine) Start(ctx context.Context) {
|
|
e.Persistence.Start(ctx)
|
|
e.Ollama.Start(ctx)
|
|
|
|
// Bootstrap local knowledge and embeddings before GPU-/graph-heavy autonomous
|
|
// workflows are allowed to run. The HTTP server is started by main in
|
|
// parallel, so this gate does not make the UI unavailable during a fresh
|
|
// bootstrap. It only prevents Security/Thinking/Autonomous Research from
|
|
// competing with the initial embedding build.
|
|
go func() {
|
|
if !e.LearningEnabled() {
|
|
e.markBootstrapReady(nil)
|
|
} else if err := e.Scan(ctx); err != nil && !errors.Is(err, ErrLearningDisabled) {
|
|
slog.Error("initial brain scan failed; autonomous workflows remain gated", "error", err)
|
|
e.markBootstrapFailure(err)
|
|
} else {
|
|
e.markBootstrapReady(nil)
|
|
}
|
|
|
|
ticker := time.NewTicker(e.Cfg.ScanInterval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
if !e.LearningEnabled() {
|
|
if !e.bootstrapIsComplete() {
|
|
e.markBootstrapReady(nil)
|
|
}
|
|
continue
|
|
}
|
|
err := e.Scan(ctx)
|
|
if err != nil && !errors.Is(err, ErrLearningDisabled) {
|
|
slog.Error("brain scan failed", "error", err)
|
|
if !e.bootstrapIsComplete() {
|
|
e.markBootstrapFailure(err)
|
|
}
|
|
continue
|
|
}
|
|
if !e.bootstrapIsComplete() {
|
|
e.markBootstrapReady(nil)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
|
|
go func() {
|
|
if e.waitBootstrap(ctx) {
|
|
e.vectorMaintenanceLoop(ctx)
|
|
}
|
|
}()
|
|
|
|
go e.enrichmentWorker(ctx)
|
|
if e.Cfg.AutoEnrich {
|
|
go e.enrichmentScheduler(ctx)
|
|
}
|
|
go e.idle(ctx)
|
|
go func() {
|
|
if !e.waitBootstrap(ctx) {
|
|
return
|
|
}
|
|
if e.GLPIKB != nil {
|
|
e.GLPIKB.Start(ctx)
|
|
}
|
|
e.startAutonomousResearch(ctx)
|
|
}()
|
|
if e.SourceInbox != nil && e.Cfg.SourceInboxEnabled {
|
|
go func() {
|
|
if e.waitBootstrap(ctx) {
|
|
e.sourceInboxLoop(ctx)
|
|
}
|
|
}()
|
|
}
|
|
if e.SourceInbox != nil {
|
|
go func() {
|
|
if e.waitBootstrap(ctx) {
|
|
e.controllerAutomationLoop(ctx)
|
|
}
|
|
}()
|
|
}
|
|
}
|
|
|
|
func (e *Engine) bootstrapIsComplete() bool {
|
|
e.stateMu.RLock()
|
|
defer e.stateMu.RUnlock()
|
|
return e.bootstrapComplete
|
|
}
|
|
|
|
func (e *Engine) markBootstrapFailure(err error) {
|
|
if err == nil || errors.Is(err, ErrLearningDisabled) {
|
|
return
|
|
}
|
|
e.stateMu.Lock()
|
|
changed := e.bootstrapError != err.Error() || e.bootstrapComplete
|
|
e.bootstrapComplete = false
|
|
e.bootstrapCompletedAt = time.Time{}
|
|
e.bootstrapError = err.Error()
|
|
e.stateMu.Unlock()
|
|
if changed && e.Broker != nil {
|
|
e.Broker.Publish(model.Activity{Type: "system.bootstrap.failed", Source: "brain", Phase: "startup", Message: "Initialer Knowledge-Scan ist fehlgeschlagen; autonome Workflows bleiben bis zu einem erfolgreichen Wiederholungsversuch gesperrt", Strength: .82, Metadata: map[string]any{"result": "blocked", "error": err.Error()}})
|
|
}
|
|
}
|
|
|
|
func (e *Engine) markBootstrapReady(err error) {
|
|
if err != nil && !errors.Is(err, ErrLearningDisabled) {
|
|
e.markBootstrapFailure(err)
|
|
return
|
|
}
|
|
e.stateMu.Lock()
|
|
wasComplete := e.bootstrapComplete
|
|
e.bootstrapComplete = true
|
|
e.bootstrapCompletedAt = time.Now().UTC()
|
|
e.bootstrapError = ""
|
|
completedAt := e.bootstrapCompletedAt
|
|
e.stateMu.Unlock()
|
|
e.bootstrapOnce.Do(func() { close(e.bootstrapReady) })
|
|
if !wasComplete && e.Broker != nil {
|
|
e.Broker.Publish(model.Activity{Type: "system.bootstrap.completed", Source: "brain", Phase: "startup", Message: "Initialer Knowledge-/Embedding-Bootstrap ist abgeschlossen; autonome Workflows werden freigegeben", Strength: .44, Metadata: map[string]any{"result": "ready", "completed_at": completedAt}})
|
|
}
|
|
}
|
|
|
|
func (e *Engine) waitBootstrap(ctx context.Context) bool {
|
|
if e.bootstrapReady == nil {
|
|
return true
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return false
|
|
case <-e.bootstrapReady:
|
|
return true
|
|
}
|
|
}
|
|
|
|
func (e *Engine) enrichmentScheduler(ctx context.Context) {
|
|
firstDelay := 12 * time.Second
|
|
if e.SpeedModeEnabled() {
|
|
firstDelay = 0
|
|
}
|
|
e.setNextEnrich(time.Now().Add(firstDelay))
|
|
timer := time.NewTimer(firstDelay)
|
|
defer timer.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-timer.C:
|
|
if !e.RequestEnrich("automatic") {
|
|
slog.Debug("automatic enrichment already queued")
|
|
}
|
|
next := time.Now().Add(e.Cfg.EnrichInterval)
|
|
e.setNextEnrich(next)
|
|
timer.Reset(e.Cfg.EnrichInterval)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (e *Engine) enrichmentWorker(ctx context.Context) {
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case trigger := <-e.enrichRequests:
|
|
if !e.waitBootstrap(ctx) {
|
|
return
|
|
}
|
|
more := e.runEnrichmentCycle(ctx, trigger)
|
|
if more && e.SpeedModeEnabled() && e.ThinkingEnabled() {
|
|
e.RequestEnrich("speed-drain")
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (e *Engine) RequestEnrich(trigger string) bool {
|
|
if !e.ThinkingEnabled() {
|
|
e.stateMu.Lock()
|
|
e.enrichResult = "disabled"
|
|
e.enrichError = ErrThinkingDisabled.Error()
|
|
e.stateMu.Unlock()
|
|
return false
|
|
}
|
|
if strings.TrimSpace(trigger) == "" {
|
|
trigger = "manual"
|
|
}
|
|
e.stateMu.Lock()
|
|
if e.enrichRunning || e.enrichResult == "queued" {
|
|
e.stateMu.Unlock()
|
|
return false
|
|
}
|
|
e.enrichResult = "queued"
|
|
e.enrichTrigger = trigger
|
|
e.enrichError = ""
|
|
e.stateMu.Unlock()
|
|
|
|
select {
|
|
case e.enrichRequests <- trigger:
|
|
e.Broker.Publish(model.Activity{Type: "think.queued", Source: "brain", Phase: "queue", Message: "AI-THINK-Zyklus wurde eingeplant", Strength: .38, Metadata: map[string]any{"trigger": trigger, "batch_size": e.Cfg.EnrichBatchSize}})
|
|
return true
|
|
default:
|
|
e.stateMu.Lock()
|
|
if e.enrichResult == "queued" && e.enrichTrigger == trigger {
|
|
e.enrichResult = "idle"
|
|
}
|
|
e.stateMu.Unlock()
|
|
return false
|
|
}
|
|
}
|
|
|
|
func (e *Engine) runEnrichmentCycle(ctx context.Context, trigger string) bool {
|
|
started := time.Now()
|
|
e.stateMu.Lock()
|
|
e.enrichRunning = true
|
|
e.enrichTrigger = trigger
|
|
e.enrichResult = "running"
|
|
e.enrichError = ""
|
|
e.lastAttempt = started.UTC()
|
|
e.enrichCycles++
|
|
e.stateMu.Unlock()
|
|
|
|
e.Broker.Publish(model.Activity{Type: "think.cycle.started", Source: "brain", Phase: "autonomous", Message: fmt.Sprintf("Autonomer AI-THINK-Zyklus startet · bis zu %d sequenzielle Prüfungen", e.Cfg.EnrichBatchSize), Strength: .72, Metadata: map[string]any{"trigger": trigger, "batch_size": e.Cfg.EnrichBatchSize, "anchors": e.Cfg.EnrichAnchors, "processing_mode": e.RuntimeSettings().ProcessingMode}})
|
|
|
|
created, rejected, checked := 0, 0, 0
|
|
exactComparisons, coarseComparisonsTotal, candidatePoolTotal := 0, 0, 0
|
|
relationsCreated, articlesCreated, articlesSkipped := 0, 0, 0
|
|
var cycleMutations graph.MutationStats
|
|
pendingArticles := make([]*pendingArticleCandidate, 0, e.Cfg.EnrichBatchSize)
|
|
result := "completed"
|
|
var cycleErr error
|
|
for step := 0; step < e.Cfg.EnrichBatchSize; step++ {
|
|
if !e.ThinkingEnabled() {
|
|
result = "disabled"
|
|
break
|
|
}
|
|
outcome, err := e.enrichOne(ctx, trigger)
|
|
if err != nil {
|
|
cycleErr = err
|
|
result = "failed"
|
|
break
|
|
}
|
|
if !outcome.Candidate {
|
|
if checked == 0 {
|
|
result = "no_candidate"
|
|
}
|
|
break
|
|
}
|
|
checked++
|
|
exactComparisons += outcome.Comparisons
|
|
coarseComparisonsTotal += outcome.CoarseComparisons
|
|
candidatePoolTotal += outcome.CandidatePool
|
|
if outcome.Created {
|
|
created++
|
|
}
|
|
if outcome.RelationCreated {
|
|
relationsCreated++
|
|
}
|
|
if outcome.ArticleCreated {
|
|
articlesCreated++
|
|
}
|
|
if outcome.ArticleSkipped {
|
|
articlesSkipped++
|
|
}
|
|
cycleMutations.Add(outcome.Mutations)
|
|
if outcome.PendingArticle != nil {
|
|
pendingArticles = append(pendingArticles, outcome.PendingArticle)
|
|
}
|
|
if outcome.Rejected {
|
|
rejected++
|
|
}
|
|
if delay := e.effectiveEnrichStepDelay(); step+1 < e.Cfg.EnrichBatchSize && delay > 0 {
|
|
select {
|
|
case <-ctx.Done():
|
|
cycleErr = ctx.Err()
|
|
result = "cancelled"
|
|
step = e.Cfg.EnrichBatchSize
|
|
case <-time.After(delay):
|
|
}
|
|
}
|
|
}
|
|
|
|
if len(pendingArticles) > 0 && cycleErr == nil {
|
|
clusterCreated, clusterSkipped := e.synthesizePendingArticleClusters(ctx, trigger, pendingArticles)
|
|
articlesCreated += clusterCreated
|
|
articlesSkipped += clusterSkipped
|
|
}
|
|
|
|
e.stateMu.Lock()
|
|
e.enrichRunning = false
|
|
e.enrichResult = result
|
|
if cycleErr != nil {
|
|
e.enrichError = cycleErr.Error()
|
|
} else {
|
|
e.enrichError = ""
|
|
}
|
|
e.enrichCreated += uint64(created)
|
|
e.enrichRejected += uint64(rejected)
|
|
e.relationsCreated += uint64(relationsCreated)
|
|
e.articlesCreated += uint64(articlesCreated)
|
|
e.articlesSkipped += uint64(articlesSkipped)
|
|
e.stateMu.Unlock()
|
|
|
|
metadata := withRunMutations(map[string]any{"trigger": trigger, "checked": checked, "created": created, "relations_created": relationsCreated, "articles_created": articlesCreated, "articles_skipped": articlesSkipped, "rejected": rejected, "duration_ms": time.Since(started).Milliseconds(), "result": result, "processing_mode": e.RuntimeSettings().ProcessingMode, "exact_comparisons": exactComparisons, "coarse_comparisons": coarseComparisonsTotal, "candidate_pool": candidatePoolTotal}, cycleMutations)
|
|
if cycleErr != nil {
|
|
e.Broker.Publish(model.Activity{Type: "think.cycle.failed", Source: "brain", Phase: "autonomous", Message: "AI-THINK-Zyklus wurde mit Fehler beendet", Strength: .45, Metadata: metadata})
|
|
slog.Warn("enrichment cycle failed", "trigger", trigger, "error", cycleErr)
|
|
return false
|
|
}
|
|
message := fmt.Sprintf("AI-THINK-Zyklus abgeschlossen · %d Relationen · %d Artikel · %d verworfen", relationsCreated, articlesCreated, rejected)
|
|
if result == "no_candidate" {
|
|
message = "AI-THINK hat im aktuell geprüften Graphbereich keinen Kandidaten oberhalb des Schwellwerts gefunden"
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "think.cycle.completed", Source: "brain", Phase: "autonomous", Message: message, Strength: .58, Metadata: metadata})
|
|
return result == "completed" && checked > 0
|
|
}
|
|
|
|
func (e *Engine) setNextEnrich(t time.Time) {
|
|
e.stateMu.Lock()
|
|
e.nextEnrich = t.UTC()
|
|
e.stateMu.Unlock()
|
|
}
|
|
|
|
func (e *Engine) setOllamaOK(ok bool) {
|
|
e.stateMu.Lock()
|
|
e.ollamaOK = ok
|
|
e.stateMu.Unlock()
|
|
}
|
|
|
|
func (e *Engine) isOllamaOK() bool {
|
|
e.stateMu.RLock()
|
|
ok := e.ollamaOK
|
|
e.stateMu.RUnlock()
|
|
return ok
|
|
}
|
|
|
|
func (e *Engine) idle(ctx context.Context) {
|
|
ticker := time.NewTicker(7 * time.Second)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
n, ok := e.Graph.IdleNode(time.Now().Unix() / 7)
|
|
if !ok {
|
|
continue
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "brain.idle", Source: "brain", Phase: "idle", Message: "Leise Hintergrundaktivität", NodeIDs: []string{n.ID}, Strength: .18})
|
|
}
|
|
}
|
|
}
|
|
func (e *Engine) configureEmbeddingDigest() int {
|
|
for _, node := range e.Ollama.NodeStatuses() {
|
|
if node.Healthy && node.Compatible && node.EmbeddingModel && strings.TrimSpace(node.EmbeddingDigest) != "" {
|
|
return e.Graph.ConfigureEmbeddingIdentity(e.Cfg.EmbeddingModel, node.EmbeddingDigest)
|
|
}
|
|
}
|
|
return 0
|
|
}
|
|
|
|
func (e *Engine) Scan(ctx context.Context) error {
|
|
if !e.LearningEnabled() {
|
|
return ErrLearningDisabled
|
|
}
|
|
e.mu.Lock()
|
|
defer e.mu.Unlock()
|
|
started := time.Now().UTC()
|
|
runID := fmt.Sprintf("learning-scan-%d", started.UnixNano())
|
|
beforeVersion := e.Graph.Version()
|
|
beforeNodes, beforeEdges, _ := e.Graph.Counts()
|
|
e.Broker.Publish(model.Activity{Type: "learning.scan.started", Source: "brain", Phase: "ingest", Message: "KB-Lernlauf gestartet: Quellen werden verglichen, Änderungen übernommen und Embeddings geprüft", Strength: .52, Metadata: map[string]any{"run_id": runID, "nodes_before": beforeNodes, "edges_before": beforeEdges, "graph_version_before": beforeVersion}})
|
|
|
|
scanResult, err := e.Scanner.ScanDetailed()
|
|
if err != nil {
|
|
meta := withRunMutations(map[string]any{"run_id": runID, "error": err.Error(), "duration_ms": time.Since(started).Milliseconds()}, graph.MutationStats{})
|
|
e.Broker.Publish(model.Activity{Type: "learning.scan.failed", Source: "brain", Phase: "ingest", Message: "KB-Lernlauf ist beim Einlesen der Wissensquellen fehlgeschlagen", Strength: .35, Metadata: meta})
|
|
return err
|
|
}
|
|
count := scanResult.Count
|
|
runMutations := scanResult.Mutations
|
|
// Repair/preserve runtime article provenance after file-owned staging nodes
|
|
// have been reconciled. This is idempotent and also heals graphs created by
|
|
// releases that accidentally deleted synthesized_from/proposes_* edges.
|
|
runMutations.Add(e.reconcileArticleProvenance())
|
|
filter := e.effectiveLearningFilter()
|
|
pendingKnowledge := len(e.Graph.KnowledgeNodesForEmbeddingScoped(filter))
|
|
fallbackVectors := e.Graph.CountVectorsByDimension(256)
|
|
|
|
// Digest changes are visible through the already running Ollama health pool;
|
|
// checking them does not require another model request. A digest change or
|
|
// remaining fallback vectors requires a one-time global repair, otherwise the
|
|
// scheduled KB scan embeds only knowledge/ai-think nodes owned by this path.
|
|
repairAllEmbeddings := false
|
|
if cleared := e.configureEmbeddingDigest(); cleared > 0 {
|
|
runMutations.VectorsDeleted += uint64(cleared)
|
|
repairAllEmbeddings = true
|
|
e.Broker.Publish(model.Activity{Type: "embedding.identity_changed", Source: "brain", Phase: "learning", Message: fmt.Sprintf("Embedding-Digest geändert · %d Vektoren werden neu gelernt", cleared), Strength: .7, Metadata: withRunMutations(map[string]any{"run_id": runID, "embedding_model": e.Cfg.EmbeddingModel, "cleared_vectors": cleared}, graph.MutationStats{VectorsDeleted: uint64(cleared)})})
|
|
}
|
|
if fallbackVectors > 0 {
|
|
repairAllEmbeddings = true
|
|
}
|
|
|
|
// The fast path is intentionally model-free: if no knowledge file changed,
|
|
// no embedding is missing and Ollama was healthy, a scheduled scan returns
|
|
// after the cheap filesystem manifest check instead of traversing/parsing the
|
|
// entire KB or pinging the model service.
|
|
needsEmbeddingWork := pendingKnowledge > 0 || repairAllEmbeddings
|
|
if !scanResult.FastPath || needsEmbeddingWork || !e.isOllamaOK() {
|
|
if !runMutations.Empty() || needsEmbeddingWork {
|
|
e.Broker.Publish(model.Activity{Type: "scan.started", Source: "brain", Phase: "ingest", Message: "Neue oder geänderte Wissenselemente werden verarbeitet", Strength: .45, Metadata: map[string]any{"run_id": runID, "pending_embeddings": pendingKnowledge, "repair_all_embeddings": repairAllEmbeddings}})
|
|
}
|
|
pingCtx, pingCancel := context.WithTimeout(ctx, 3*time.Second)
|
|
pingErr := e.Ollama.Ping(pingCtx)
|
|
pingCancel()
|
|
if pingErr != nil {
|
|
slog.Warn("Ollama unavailable; using deterministic local fallback", "error", pingErr)
|
|
fallbackStats := e.ensureFallbackEmbeddings(false)
|
|
runMutations.Add(fallbackStats)
|
|
e.setOllamaOK(false)
|
|
} else {
|
|
// Local fallback vectors use 256 dimensions. Once Ollama is healthy,
|
|
// replace all placeholders, including those created by other workflows.
|
|
if cleared := e.Graph.ClearVectorsByDimension(256); cleared > 0 {
|
|
runMutations.VectorsDeleted += uint64(cleared)
|
|
repairAllEmbeddings = true
|
|
}
|
|
embedStats, embedErr := e.ensureEmbeddings(ctx, repairAllEmbeddings)
|
|
runMutations.Add(embedStats)
|
|
if embedErr != nil {
|
|
slog.Warn("Ollama embeddings failed; using deterministic local fallback", "error", embedErr)
|
|
fallbackStats := e.ensureFallbackEmbeddings(repairAllEmbeddings)
|
|
runMutations.Add(fallbackStats)
|
|
e.setOllamaOK(false)
|
|
} else {
|
|
e.setOllamaOK(true)
|
|
}
|
|
}
|
|
}
|
|
|
|
// External evidence belongs to Research/Security workflows and is therefore
|
|
// not charged to the learning run. If one of those workflows previously
|
|
// failed to embed its node, repair it here as a separate explicitly
|
|
// attributed embedding workflow so readiness can self-heal without lying
|
|
// about Learning costs.
|
|
if e.isOllamaOK() {
|
|
e.repairMissingExternalEmbeddings(ctx)
|
|
}
|
|
|
|
// Build a sparse semantic Knowledge<->Knowledge layer from already existing
|
|
// embeddings. The calculation itself needs no model call and may either run
|
|
// locally on the Brain CPU or be claimed by a compute-capable Source Agent.
|
|
// Only 256-dimensional deterministic fallback vectors are excluded because
|
|
// mixing them with the configured embedding space would make thresholds lie.
|
|
// Scheduled reevaluation/layout maintenance has its own loop and therefore
|
|
// continues even while Learning is paused. Learning scans rebuild the vector
|
|
// layer only when their own ingest/embedding work changed the knowledge space
|
|
// or when no vector layer exists yet.
|
|
vectorLayerNeeded := !scanResult.FastPath || needsEmbeddingWork || !e.Graph.HasEdgesByOrigin(graph.VectorMathOrigin)
|
|
vectorReady := e.Graph.CountVectorsByDimension(256) == 0
|
|
if e.Cfg.VectorGraphEnabled && vectorReady && vectorLayerNeeded {
|
|
startedVectorGraph := time.Now()
|
|
execution, vectorErr := e.rebuildVectorSemanticLayer(ctx, runID, filter)
|
|
if vectorErr != nil {
|
|
slog.Warn("vector graph rebuild skipped", "error", vectorErr)
|
|
e.Broker.Publish(model.Activity{Type: "vector.graph.failed", Source: "brain", Phase: "semantic-linking", Message: "Mathematische Vektorverknüpfung konnte nicht sicher aktualisiert werden", Strength: .34, Metadata: map[string]any{"run_id": runID, "error": vectorErr.Error(), "agent_required": e.Cfg.VectorGraphAgentRequired}})
|
|
} else {
|
|
stats, vectorMutations := execution.Stats, execution.Mutations
|
|
runMutations.Add(vectorMutations)
|
|
e.stateMu.Lock()
|
|
e.lastVectorGraph = time.Now().UTC()
|
|
e.stateMu.Unlock()
|
|
e.Broker.Publish(model.Activity{Type: "vector.graph.rebuilt", Source: "brain", Phase: "semantic-linking", Message: fmt.Sprintf("Mathematische Vektorverknüpfung: %d Primär- + %d Orphan-Kanten aus %d Embeddings", stats.Links, stats.OrphanLinks, stats.Indexed), Strength: .62, Metadata: withRunMutations(map[string]any{
|
|
"run_id": runID, "algorithm": "mutual-knn-local-scaling-v1", "orphan_algorithm": "orphan-knn-local-scaling-v1", "no_model_call": true,
|
|
"indexed": stats.Indexed, "links": stats.Links, "reciprocal_links": stats.ReciprocalLinks,
|
|
"candidate_pairs": stats.CandidatePairs, "exact_comparisons": stats.ExactComparisons,
|
|
"orphan_pass_enabled": e.Cfg.VectorGraphOrphanPass, "orphan_focus": stats.OrphanFocus, "orphan_links": stats.OrphanLinks,
|
|
"orphan_exact_comparisons": stats.OrphanStats.ExactComparisons, "orphan_candidate_pairs": stats.OrphanStats.CandidatePairs,
|
|
"position_updates": stats.PositionUpdates, "layout_enabled": e.Cfg.VectorGraphLayout, "layout_relaxation_configured": e.Cfg.VectorGraphRelaxLayout, "layout_mode": execution.LayoutMode, "layout_applied": execution.LayoutDue && stats.PositionUpdates > 0, "periodic_refresh": false, "maintenance_owned": true, "reevaluate_interval": e.Cfg.VectorGraphReevaluateInterval.String(),
|
|
"agent_offloaded": execution.Offloaded, "agent_id": execution.AgentID, "agent_compute_ms": execution.ComputeMS, "agent_fallback_reason": execution.FallbackReason,
|
|
"speed_mode": e.SpeedModeEnabled(), "cpu_workers": e.vectorPrimaryConfig(execution.LayoutDue).Workers,
|
|
"duration_ms": time.Since(startedVectorGraph).Milliseconds(),
|
|
}, vectorMutations)})
|
|
}
|
|
}
|
|
|
|
e.stateMu.Lock()
|
|
e.lastScan = time.Now().UTC()
|
|
e.stateMu.Unlock()
|
|
nodes, edges, version := e.Graph.Counts()
|
|
if !runMutations.Empty() {
|
|
e.Broker.Publish(model.Activity{Type: "graph.updated", Source: "brain", Phase: "indexed", Message: fmt.Sprintf("%d Wissenselemente · %d Nodes · %d Edges", count, nodes, edges), Strength: .55, Metadata: withRunMutations(map[string]any{"run_id": runID, "nodes": nodes, "edges": edges, "knowledge_elements": count}, runMutations)})
|
|
}
|
|
result := "unchanged"
|
|
if !runMutations.Empty() {
|
|
result = "updated"
|
|
}
|
|
meta := map[string]any{
|
|
"run_id": runID, "result": result, "duration_ms": time.Since(started).Milliseconds(), "knowledge_elements": count,
|
|
"nodes_before": beforeNodes, "nodes_after": nodes, "edges_before": beforeEdges, "edges_after": edges,
|
|
"graph_version_before": beforeVersion, "graph_version_after": version,
|
|
"nodes_created": runMutations.NodesCreated, "nodes_updated": runMutations.NodesUpdated, "nodes_deleted": runMutations.NodesDeleted,
|
|
"edges_created": runMutations.EdgesCreated, "edges_updated": runMutations.EdgesUpdated, "edges_deleted": runMutations.EdgesDeleted,
|
|
"vectors_created": runMutations.VectorsCreated, "vectors_updated": runMutations.VectorsUpdated, "vectors_deleted": runMutations.VectorsDeleted,
|
|
"pending_embeddings": len(e.Graph.KnowledgeNodesForEmbeddingScoped(filter)), "ollama_ok": e.isOllamaOK(),
|
|
"manifest_fast_path": scanResult.FastPath, "manifest_changed": scanResult.ManifestChanged, "full_content_verify": scanResult.FullVerify,
|
|
}
|
|
meta = withRunMutations(meta, runMutations)
|
|
e.Broker.Publish(model.Activity{Type: "learning.scan.completed", Source: "brain", Phase: "indexed", Message: fmt.Sprintf("KB-Lernlauf abgeschlossen · %d Wissenselemente · %d neue Nodes · %d neue Edges · %d neue/neu berechnete Embeddings", count, runMutations.NodesCreated, runMutations.EdgesCreated, runMutations.VectorsCreated+runMutations.VectorsUpdated), Strength: .64, Metadata: meta})
|
|
return nil
|
|
}
|
|
|
|
func (e *Engine) repairMissingExternalEmbeddings(ctx context.Context) {
|
|
all := e.Graph.NodesForEmbeddingScoped(graph.NodeFilter{})
|
|
pending := make([]model.Node, 0)
|
|
for _, node := range all {
|
|
if node.Kind == "external" {
|
|
pending = append(pending, node)
|
|
}
|
|
}
|
|
if len(pending) == 0 || e.Ollama == nil {
|
|
return
|
|
}
|
|
started := time.Now()
|
|
var total graph.MutationStats
|
|
for start := 0; start < len(pending); start += 16 {
|
|
end := start + 16
|
|
if end > len(pending) {
|
|
end = len(pending)
|
|
}
|
|
texts := make([]string, 0, end-start)
|
|
for _, node := range pending[start:end] {
|
|
texts = append(texts, embeddingText(node))
|
|
}
|
|
vectors, err := e.Ollama.Embed(ctx, texts)
|
|
if err != nil || len(vectors) != len(texts) {
|
|
e.Broker.Publish(model.Activity{Type: "embedding.external_repair.failed", Source: "ollama", Phase: "embedding", Message: "Fehlende externe Embeddings konnten noch nicht repariert werden", Strength: .34, Metadata: map[string]any{"pending_external": len(pending), "error": errorString(err), "duration_ms": time.Since(started).Milliseconds()}})
|
|
return
|
|
}
|
|
for i, vector := range vectors {
|
|
if len(vector) == 0 {
|
|
continue
|
|
}
|
|
total.Add(e.Graph.SetVectorWithStats(pending[start+i].ID, vector))
|
|
}
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "embedding.external_repair.completed", Source: "ollama", Phase: "embedding", NodeIDs: nodeIDsFromNodes(pending), Message: fmt.Sprintf("%d fehlende externe Embeddings wurden repariert", total.VectorsCreated+total.VectorsUpdated), Strength: .54, Metadata: withRunMutations(map[string]any{"pending_external": len(pending), "duration_ms": time.Since(started).Milliseconds(), "model": e.Cfg.EmbeddingModel}, total)})
|
|
}
|
|
|
|
func (e *Engine) embeddingPending(includeExternal bool) []model.Node {
|
|
if includeExternal {
|
|
return e.Graph.NodesForEmbeddingScoped(e.effectiveLearningFilter())
|
|
}
|
|
return e.Graph.KnowledgeNodesForEmbeddingScoped(e.effectiveLearningFilter())
|
|
}
|
|
|
|
func (e *Engine) ensureEmbeddings(ctx context.Context, includeExternal bool) (graph.MutationStats, error) {
|
|
pending := e.embeddingPending(includeExternal)
|
|
var total graph.MutationStats
|
|
if len(pending) == 0 {
|
|
return total, nil
|
|
}
|
|
for start := 0; start < len(pending); start += 16 {
|
|
end := start + 16
|
|
if end > len(pending) {
|
|
end = len(pending)
|
|
}
|
|
texts := make([]string, 0, end-start)
|
|
for _, n := range pending[start:end] {
|
|
texts = append(texts, embeddingText(n))
|
|
}
|
|
batchStarted := time.Now()
|
|
cctx, cancel := context.WithTimeout(ctx, 4*time.Minute)
|
|
vecs, err := e.Ollama.Embed(cctx, texts)
|
|
cancel()
|
|
if err != nil {
|
|
return total, err
|
|
}
|
|
if len(vecs) != len(texts) {
|
|
return total, fmt.Errorf("embedding response count mismatch: got %d vectors for %d inputs", len(vecs), len(texts))
|
|
}
|
|
var batchStats graph.MutationStats
|
|
for i, v := range vecs {
|
|
if len(v) == 0 {
|
|
return total, fmt.Errorf("embedding response %d is empty", start+i)
|
|
}
|
|
batchStats.Add(e.Graph.SetVectorWithStats(pending[start+i].ID, v))
|
|
}
|
|
total.Add(batchStats)
|
|
ids := []string{}
|
|
for _, n := range pending[start:end] {
|
|
ids = append(ids, n.ID)
|
|
}
|
|
meta := withRunMutations(map[string]any{"batch_count": len(ids), "model": e.Cfg.EmbeddingModel, "batch_start": start, "batch_total": len(pending), "scope": map[bool]string{true: "all", false: "knowledge"}[includeExternal], "duration_ms": time.Since(batchStarted).Milliseconds()}, batchStats)
|
|
e.Broker.Publish(model.Activity{Type: "embedding.batch", Source: "ollama", Phase: "embedding", Message: fmt.Sprintf("EmbeddingGemma verarbeitet %d Elemente", len(ids)), NodeIDs: ids, Strength: .38, Metadata: meta})
|
|
}
|
|
return total, nil
|
|
}
|
|
|
|
func (e *Engine) ensureFallbackEmbeddings(includeExternal bool) graph.MutationStats {
|
|
var total graph.MutationStats
|
|
for _, n := range e.embeddingPending(includeExternal) {
|
|
total.Add(e.Graph.SetVectorWithStats(n.ID, hashEmbedding(embeddingText(n), 256)))
|
|
}
|
|
return total
|
|
}
|
|
|
|
func embeddingText(n model.Node) string {
|
|
return strings.TrimSpace(n.Label + "\n" + strings.Join(n.Categories, " · ") + "\n" + strings.Join(n.Keywords, " · ") + "\n" + n.Summary)
|
|
}
|
|
func hashEmbedding(s string, dims int) []float64 {
|
|
v := make([]float64, dims)
|
|
tokens := strings.FieldsFunc(strings.ToLower(s), func(r rune) bool { return !unicode.IsLetter(r) && !unicode.IsDigit(r) })
|
|
for _, t := range tokens {
|
|
if t == "" {
|
|
continue
|
|
}
|
|
h := sha256.Sum256([]byte(t))
|
|
idx := (int(h[0])<<8 | int(h[1])) % dims
|
|
sign := 1.0
|
|
if h[2]&1 == 1 {
|
|
sign = -1
|
|
}
|
|
v[idx] += sign * (1 + float64(h[3])/255)
|
|
}
|
|
var norm float64
|
|
for _, x := range v {
|
|
norm += x * x
|
|
}
|
|
if norm > 0 {
|
|
norm = math.Sqrt(norm)
|
|
for i := range v {
|
|
v[i] /= norm
|
|
}
|
|
}
|
|
return v
|
|
}
|
|
|
|
func (e *Engine) similarKnowledge(query []float64, limit int, filter graph.NodeFilter, maxAIDepth int) ([]model.Hit, graph.ClusterSearchStats) {
|
|
if e.RuntimeSettings().ProcessingMode != "clustered" {
|
|
return e.Graph.SimilarFiltered(query, limit, filter), graph.ClusterSearchStats{}
|
|
}
|
|
candidateLimit := e.Cfg.ClusterCandidatesPerAnchor
|
|
if candidateLimit < limit*12 {
|
|
candidateLimit = limit * 12
|
|
}
|
|
return e.Graph.SimilarClusteredFiltered(query, limit, candidateLimit, filter, maxAIDepth, e.Cfg.ClusterHashBits, e.Cfg.ClusterHashTables)
|
|
}
|
|
|
|
func (e *Engine) Query(ctx context.Context, q string) (model.QueryResponse, error) {
|
|
e.interactiveInflight.Add(1)
|
|
defer e.interactiveInflight.Add(-1)
|
|
start := time.Now()
|
|
q = strings.TrimSpace(q)
|
|
if len([]rune(q)) < 2 {
|
|
return model.QueryResponse{}, fmt.Errorf("query is too short")
|
|
}
|
|
queryRunID := graph.ID("query-run", q, fmt.Sprintf("%d", start.UnixNano()))
|
|
queryMeta := func(extra map[string]any) map[string]any {
|
|
meta := map[string]any{"run_id": queryRunID}
|
|
for key, value := range extra {
|
|
meta[key] = value
|
|
}
|
|
return meta
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "query.started", Source: "ui", Phase: "perception", Query: q, Message: "Anfrage trifft im neuronalen Feld ein", Strength: 1, Metadata: queryMeta(nil)})
|
|
vecs, err := e.Ollama.Embed(ctx, []string{q})
|
|
if err != nil || len(vecs) == 0 {
|
|
vecs = [][]float64{hashEmbedding(q, 256)}
|
|
}
|
|
hits, retrievalStats := e.similarKnowledge(vecs[0], e.Cfg.TopK, e.effectiveLearningFilter(), 0)
|
|
if e.RuntimeSettings().ProcessingMode == "clustered" {
|
|
e.Broker.Publish(model.Activity{Type: "query.retrieval.clustered", Source: "brain", Phase: "retrieval", Query: q, Message: fmt.Sprintf("Cluster-Retrieval: %d exakte Cosine-Prüfungen nach %d Hash-Vergleichen", retrievalStats.ExactComparisons, retrievalStats.CoarseComparisons), Strength: .28, Metadata: queryMeta(map[string]any{"processing_mode": "clustered", "indexed_nodes": retrievalStats.IndexedNodes, "coarse_comparisons": retrievalStats.CoarseComparisons, "exact_comparisons": retrievalStats.ExactComparisons, "candidate_pool": retrievalStats.CandidatePool})})
|
|
}
|
|
nodeIDs := make([]string, 0, len(hits))
|
|
for i, h := range hits {
|
|
nodeIDs = append(nodeIDs, h.NodeID)
|
|
e.Broker.Publish(model.Activity{Type: "node.activated", Source: "brain", Phase: "retrieval", Query: q, NodeIDs: []string{h.NodeID}, Message: fmt.Sprintf("Treffer %d · %.0f%% · %s", i+1, h.Score*100, h.Label), Strength: math.Max(.25, h.Score), Metadata: queryMeta(nil)})
|
|
if !e.SpeedModeEnabled() {
|
|
time.Sleep(55 * time.Millisecond)
|
|
}
|
|
}
|
|
edgeIDs := e.Graph.ConnectingEdges(nodeIDs)
|
|
if len(edgeIDs) > 0 {
|
|
e.Broker.Publish(model.Activity{Type: "edges.traversed", Source: "brain", Phase: "association", Query: q, NodeIDs: nodeIDs, EdgeIDs: edgeIDs, Message: fmt.Sprintf("%d Wissensverbindungen werden durchlaufen", len(edgeIDs)), Strength: .92, Metadata: queryMeta(nil)})
|
|
}
|
|
answer := e.fallbackAnswer(q, hits)
|
|
used := append([]string(nil), nodeIDs...)
|
|
var uncertainties []string
|
|
if e.isOllamaOK() && len(hits) > 0 {
|
|
system := "Du beantwortest Fragen ausschließlich aus dem bereitgestellten Wissensgraphen. Markiere Unklarheiten offen. Gib valides JSON nach Schema zurück. used_node_ids dürfen nur IDs aus dem Kontext sein."
|
|
user := e.answerContext(q, hits)
|
|
var dec model.AnswerDecision
|
|
if err := e.Ollama.ChatJSON(ctx, system, user, answerSchema(), &dec); err == nil && strings.TrimSpace(dec.Answer) != "" {
|
|
answer = dec.Answer
|
|
used = validIDs(dec.UsedNodeIDs, nodeIDs)
|
|
uncertainties = dec.Uncertainties
|
|
} else if err != nil {
|
|
slog.Warn("structured answer failed; fallback used", "error", err)
|
|
}
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "query.completed", Source: "brain", Phase: "synthesis", Query: q, NodeIDs: used, EdgeIDs: e.Graph.ConnectingEdges(used), Message: "Antwortsynthese abgeschlossen", Strength: 1, Metadata: queryMeta(map[string]any{"duration_ms": time.Since(start).Milliseconds(), "hit_count": len(hits), "used_nodes": len(used), "uncertainty_count": len(uncertainties)})})
|
|
response := model.QueryResponse{Query: q, Answer: answer, Hits: hits, UsedNodeIDs: used, Uncertainties: uncertainties, DurationMS: time.Since(start).Milliseconds()}
|
|
if e.Cfg.AutonomousResearchQueryTriggers && e.AutonomousResearchEnabled() && (len(hits) == 0 || len(uncertainties) > 0) {
|
|
questions := append([]string(nil), uncertainties...)
|
|
if len(questions) == 0 {
|
|
questions = []string{q}
|
|
}
|
|
priority := .74
|
|
if len(hits) == 0 {
|
|
priority = .88
|
|
}
|
|
go func(request model.ResearchTaskRequest) {
|
|
queueCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if _, _, err := e.QueueResearchTask(queueCtx, request); err != nil {
|
|
slog.Debug("query uncertainty could not be queued for autonomous research", "error", err)
|
|
}
|
|
}(model.ResearchTaskRequest{Topic: q, Questions: questions, SeedNodeIDs: used, Priority: priority, RequestedBy: "query", Reason: "knowledge_answer_insufficient", Metadata: map[string]any{"hit_count": len(hits), "uncertainty_count": len(uncertainties)}})
|
|
}
|
|
return response, nil
|
|
}
|
|
func (e *Engine) answerContext(q string, hits []model.Hit) string {
|
|
var b strings.Builder
|
|
b.WriteString("FRAGE:\n" + q + "\n\nKONTEXT:\n")
|
|
used := 0
|
|
for _, h := range hits {
|
|
n, ok := e.Graph.GetNode(h.NodeID)
|
|
if !ok {
|
|
continue
|
|
}
|
|
part := fmt.Sprintf("\nNODE_ID: %s\nTITEL: %s\nSTATUS: %s\nKATEGORIEN: %s\nINHALT: %s\n", n.ID, n.Label, n.Status, strings.Join(n.Categories, ", "), n.Summary)
|
|
if used+len(part) > e.Cfg.MaxContextChars {
|
|
break
|
|
}
|
|
b.WriteString(part)
|
|
used += len(part)
|
|
}
|
|
return b.String()
|
|
}
|
|
func (e *Engine) fallbackAnswer(q string, hits []model.Hit) string {
|
|
if len(hits) == 0 {
|
|
return "Im aktuellen Wissensgraphen wurde kein belastbarer Zusammenhang gefunden."
|
|
}
|
|
var b strings.Builder
|
|
b.WriteString("Die stärksten passenden Wissensbereiche sind: ")
|
|
for i, h := range hits {
|
|
if i >= 4 {
|
|
break
|
|
}
|
|
if i > 0 {
|
|
b.WriteString("; ")
|
|
}
|
|
b.WriteString(h.Label)
|
|
}
|
|
b.WriteString(". Die Visualisierung zeigt die zugehörigen Aktivierungspfade. Ohne erreichbares Qwen-Modell bleibt dies eine Retrieval-Zusammenfassung.")
|
|
return b.String()
|
|
}
|
|
|
|
func (e *Engine) EnrichOne(ctx context.Context) error {
|
|
if !e.ThinkingEnabled() {
|
|
return ErrThinkingDisabled
|
|
}
|
|
outcome, err := e.enrichOne(ctx, "direct")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
e.stateMu.Lock()
|
|
if outcome.RelationCreated {
|
|
e.relationsCreated++
|
|
e.enrichCreated++
|
|
}
|
|
if outcome.ArticleCreated {
|
|
e.articlesCreated++
|
|
}
|
|
if outcome.ArticleSkipped {
|
|
e.articlesSkipped++
|
|
}
|
|
if outcome.Rejected {
|
|
e.enrichRejected++
|
|
}
|
|
e.stateMu.Unlock()
|
|
return nil
|
|
}
|
|
|
|
func (e *Engine) enrichOne(ctx context.Context, trigger string) (EnrichOutcome, error) {
|
|
if !e.ThinkingEnabled() {
|
|
return EnrichOutcome{Result: "disabled"}, ErrThinkingDisabled
|
|
}
|
|
e.mu.Lock()
|
|
defer e.mu.Unlock()
|
|
|
|
if !e.isOllamaOK() {
|
|
pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
|
err := e.Ollama.Ping(pingCtx)
|
|
cancel()
|
|
if err != nil {
|
|
e.Broker.Publish(model.Activity{Type: "think.paused", Source: "brain", Phase: "waiting", Message: "AI-THINK wartet auf ein erreichbares Ollama/Qwen-Modell", Strength: .25, Metadata: map[string]any{"trigger": trigger, "error": err.Error()}})
|
|
return EnrichOutcome{Result: "ollama_unavailable"}, fmt.Errorf("Ollama/Qwen is unavailable; no AI edge or AI-THINK draft was created: %w", err)
|
|
}
|
|
e.setOllamaOK(true)
|
|
}
|
|
|
|
processingMode := e.RuntimeSettings().ProcessingMode
|
|
var a, b model.Node
|
|
var sim float64
|
|
var ok bool
|
|
comparisons, coarseComparisons, indexedNodes, candidatePool := 0, 0, 0, 0
|
|
candidateSource := "embedding_search"
|
|
if e.Cfg.ThinkingVectorGuided && e.Graph.HasEdgesByOrigin(graph.VectorMathOrigin) {
|
|
var vectorStats graph.VectorNeighborCandidateStats
|
|
a, b, sim, ok, vectorStats = e.Graph.NextVectorNeighborPairScoped(e.effectiveThinkingFilter(), e.Cfg.ArticleMaxGenerationDepth)
|
|
if ok {
|
|
candidateSource = "vector_graph"
|
|
candidatePool = vectorStats.Candidates
|
|
indexedNodes = vectorStats.Candidates
|
|
}
|
|
}
|
|
if !ok && processingMode == "clustered" {
|
|
var stats graph.ClusterSearchStats
|
|
a, b, sim, ok, stats = e.Graph.NextPairClusteredScopedDepth(e.Cfg.SimilarityThreshold, e.Cfg.EnrichAnchors, e.effectiveThinkingFilter(), e.Cfg.ArticleMaxGenerationDepth, e.Cfg.ClusterHashBits, e.Cfg.ClusterHashTables, e.Cfg.ClusterCandidatesPerAnchor)
|
|
comparisons = stats.ExactComparisons
|
|
coarseComparisons = stats.CoarseComparisons
|
|
indexedNodes = stats.IndexedNodes
|
|
candidatePool = stats.CandidatePool
|
|
} else if !ok {
|
|
a, b, sim, ok, comparisons = e.Graph.NextPairScopedDepth(e.Cfg.SimilarityThreshold, e.Cfg.EnrichAnchors, e.effectiveThinkingFilter(), e.Cfg.ArticleMaxGenerationDepth)
|
|
}
|
|
if !ok {
|
|
e.stateMu.Lock()
|
|
e.lastAttempt = time.Now().UTC()
|
|
e.stateMu.Unlock()
|
|
e.Broker.Publish(model.Activity{Type: "think.no_candidate", Source: "brain", Phase: "candidate-search", Message: "Im aktuell geprüften Graphbereich wurde keine ungeprüfte Beziehung oberhalb des Ähnlichkeitsschwellwerts gefunden", Strength: .28, Metadata: map[string]any{"trigger": trigger, "threshold": e.Cfg.SimilarityThreshold, "anchors": e.Cfg.EnrichAnchors, "comparisons": comparisons, "exact_comparisons": comparisons, "coarse_comparisons": coarseComparisons, "indexed_nodes": indexedNodes, "candidate_pool": candidatePool, "processing_mode": processingMode, "candidate_source": candidateSource}})
|
|
return EnrichOutcome{Result: "no_candidate", Comparisons: comparisons, CoarseComparisons: coarseComparisons, IndexedNodes: indexedNodes, CandidatePool: candidatePool}, nil
|
|
}
|
|
|
|
now := time.Now().UTC()
|
|
e.stateMu.Lock()
|
|
e.lastAttempt = now
|
|
e.lastEnrich = now
|
|
e.stateMu.Unlock()
|
|
e.Broker.Publish(model.Activity{Type: "think.started", Source: "brain", Phase: "association", NodeIDs: []string{a.ID, b.ID}, Message: fmt.Sprintf("Verwandtschaft wird geprüft · %.0f%% semantische Nähe", sim*100), Strength: .88, Metadata: map[string]any{"trigger": trigger, "semantic_similarity": sim, "source_label": a.Label, "target_label": b.Label, "model": e.Cfg.ChatModel, "candidate_comparisons": comparisons, "exact_comparisons": comparisons, "coarse_comparisons": coarseComparisons, "indexed_nodes": indexedNodes, "candidate_pool": candidatePool, "processing_mode": processingMode, "candidate_source": candidateSource}})
|
|
|
|
system := "Du führst ausschließlich eine Relationserkennung für einen Wissensgraphen durch. Analysiere zwei interne Wissenseinträge, erfinde keine Fakten und entscheide, ob eine belastbare Beziehung besteht. Schreibe keinen Artikel und keine technische Synthese. Bei sehr hoher semantischer Nähe und same_topic/related_to ist normalerweise keine Webrecherche nötig. Wenn externe Fakten für eine kausale, abhängige, widersprüchliche oder zeitkritische Relationsentscheidung fehlen, setze needs_research=true. research_query darf ausschließlich fachliche Begriffe und sichtbare Titel enthalten, niemals interne Node-IDs, Hashes oder Datenbankkennungen. Gib ausschließlich JSON nach Schema zurück."
|
|
var decision model.RelationDecision
|
|
if err := e.Ollama.ChatJSON(ctx, system, relationContext(a, b, sim), relationSchema(), &decision); err != nil {
|
|
e.Broker.Publish(model.Activity{Type: "think.failed", Source: "brain", Phase: "inference", NodeIDs: []string{a.ID, b.ID}, Message: "Qwen-Beziehungsanalyse ist fehlgeschlagen; es wurde nichts gespeichert", Strength: .35, Metadata: map[string]any{"trigger": trigger, "error": err.Error(), "model": e.Cfg.ChatModel}})
|
|
return EnrichOutcome{Result: "inference_failed", Candidate: true, Comparisons: comparisons, CoarseComparisons: coarseComparisons, IndexedNodes: indexedNodes, CandidatePool: candidatePool}, fmt.Errorf("relation inference failed: %w", err)
|
|
}
|
|
|
|
var researchResults []model.ResearchResult
|
|
if decision.NeedsResearch && !e.relationResearchNeeded(a, b, sim, decision) {
|
|
if e.Broker != nil {
|
|
e.Broker.Publish(model.Activity{Type: "think.research.skipped", Source: "brain", Phase: "research-routing", NodeIDs: []string{a.ID, b.ID}, Message: "Externe Relationsrecherche wurde übersprungen; die interne Same-Topic-Beziehung ist bereits ausreichend belegt", Strength: .32, Metadata: map[string]any{"trigger": trigger, "semantic_similarity": sim, "relation_type": safeRelation(decision.RelationType), "confidence": decision.Confidence, "reason": "high_similarity_internal_relation"}})
|
|
}
|
|
decision.NeedsResearch = false
|
|
decision.ResearchQuery = ""
|
|
}
|
|
if decision.NeedsResearch && e.ResearchEnabledForRuntime() && strings.TrimSpace(decision.ResearchQuery) != "" {
|
|
researchID := newResearchRunID("relation-research-v10", decision.ResearchQuery)
|
|
researchStarted := time.Now()
|
|
e.Broker.Publish(model.Activity{Type: "research.started", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "Unklarheit erkannt · Relationsrecherche prüft externe Volltextbelege einzeln", Strength: .9, Metadata: map[string]any{"trigger": trigger, "research_id": researchID, "research_query": decision.ResearchQuery, "source_label": a.Label, "target_label": b.Label, "evidence_gate": "relation-fulltext-v10", "animation_min_ms": 2000}})
|
|
|
|
acceptedEvidence, researchMetadata, researchErr := e.collectRelationResearchEvidence(ctx, trigger, a, b, decision)
|
|
researchMetadata["research_id"] = researchID
|
|
researchMetadata["duration_ms"] = time.Since(researchStarted).Milliseconds()
|
|
if researchErr != nil {
|
|
researchMetadata["error"] = researchErr.Error()
|
|
decision.Related = false
|
|
decision.Confidence = math.Min(decision.Confidence, math.Max(0, e.Cfg.RelationThreshold-.01))
|
|
decision.Explanation = strings.TrimSpace(decision.Explanation + " · erforderliche externe Relationsrecherche fehlgeschlagen")
|
|
e.Broker.Publish(model.Activity{Type: "research.failed", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "Relationsrecherche ist fehlgeschlagen; die beweisabhängige Relation wird nicht übernommen", Strength: .35, Metadata: researchMetadata})
|
|
} else {
|
|
researchResults = acceptedEvidence
|
|
e.Broker.Publish(model.Activity{Type: "research.results", Source: "brain", Phase: "research-results", NodeIDs: []string{a.ID, b.ID}, Message: relationResearchResultMessage(researchMetadata), Strength: .82, Metadata: researchMetadata})
|
|
if len(acceptedEvidence) == 0 {
|
|
decision.Related = false
|
|
decision.Confidence = math.Min(decision.Confidence, math.Max(0, e.Cfg.RelationThreshold-.01))
|
|
decision.Explanation = strings.TrimSpace(decision.Explanation + " · erforderliche externe Relationsrecherche lieferte keinen belastbaren Volltextbeleg")
|
|
e.Broker.Publish(model.Activity{Type: "think.research.insufficient", Source: "brain", Phase: "research-routing", NodeIDs: []string{a.ID, b.ID}, Message: "Die Relation benötigte externe Fakten, aber kein Suchtreffer bestand die Volltext-Evidenzprüfung", Strength: .4, Metadata: researchMetadata})
|
|
} else {
|
|
var reviewed model.RelationDecision
|
|
reviewSystem := "Bewerte die Beziehung erneut anhand der zwei internen Wissenseinträge und ausschließlich der beigefügten, einzeln geprüften Volltextbelege. Die Belege haben bereits Topic-, Quellenqualitäts- und Volltext-Relevanz-Gates bestanden. Erfinde nichts. Interne Node-IDs sind keine Websuchbegriffe. Wenn die Relation weiterhin nicht belastbar ist, setze related=false. Gib ausschließlich JSON nach Schema zurück."
|
|
if err := e.Ollama.ChatJSON(ctx, reviewSystem, relationContextWithResearch(a, b, sim, acceptedEvidence), relationSchema(), &reviewed); err != nil {
|
|
slog.Warn("relation evidence review failed; rejecting evidence-dependent relation", "error", err)
|
|
researchResults = nil
|
|
decision.Related = false
|
|
decision.Confidence = math.Min(decision.Confidence, math.Max(0, e.Cfg.RelationThreshold-.01))
|
|
decision.Explanation = strings.TrimSpace(decision.Explanation + " · zweiter Relationsreview mit geprüfter Evidenz fehlgeschlagen")
|
|
} else {
|
|
decision = reviewed
|
|
}
|
|
}
|
|
researchMetadata["result"] = "completed"
|
|
e.Broker.Publish(model.Activity{Type: "research.completed", Source: "brain", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: relationResearchResultMessage(researchMetadata), Strength: .72, Metadata: researchMetadata})
|
|
}
|
|
}
|
|
|
|
status := "staging"
|
|
if !decision.Related || decision.Confidence < e.Cfg.RelationThreshold {
|
|
status = "rejected"
|
|
}
|
|
edge := model.Edge{
|
|
Source: a.ID, Target: b.ID, Type: safeRelation(decision.RelationType), Origin: "ai-inference", Status: status,
|
|
Confidence: decision.Confidence, Weight: math.Max(.2, decision.Confidence), Explanation: decision.Explanation,
|
|
Evidence: []model.Evidence{{NodeID: a.ID, URI: a.URI, Excerpt: clamp(a.Summary, 220)}, {NodeID: b.ID, URI: b.URI, Excerpt: clamp(b.Summary, 220)}},
|
|
Metadata: map[string]any{"semantic_similarity": sim, "model": e.Cfg.ChatModel, "research_result_count": len(researchResults), "trigger": trigger, "candidate_source": candidateSource},
|
|
}
|
|
edgeMutations := e.Graph.UpsertEdgeWithStats(edge)
|
|
edge.ID = graph.EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin)
|
|
outcome := EnrichOutcome{Result: status, Candidate: true, Comparisons: comparisons, CoarseComparisons: coarseComparisons, IndexedNodes: indexedNodes, CandidatePool: candidatePool, Mutations: edgeMutations}
|
|
if status == "staging" {
|
|
outcome.Created = true
|
|
outcome.RelationCreated = true
|
|
if len(researchResults) > 0 {
|
|
refs := e.addResearch(ctx, a, b, researchResults)
|
|
outcome.Mutations.Add(refs.Mutations)
|
|
e.Broker.Publish(model.Activity{Type: "research.ingested", Source: "searxng", Phase: "research-ingest", NodeIDs: append([]string{a.ID, b.ID}, refs.NodeIDs...), EdgeIDs: refs.EdgeIDs, Message: fmt.Sprintf("%d Webquellen wurden nach akzeptierter Relation in den Graphen übernommen", len(refs.NodeIDs)), Strength: .82, Metadata: map[string]any{"trigger": trigger, "result_node_ids": refs.NodeIDs, "result_edge_ids": refs.EdgeIDs, "materialization": "accepted_relation_only"}})
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "think.relation.created", Source: "brain", Phase: "relation", NodeIDs: []string{a.ID, b.ID}, EdgeIDs: []string{edge.ID}, Message: "Belastbare Wissensrelation wurde als überprüfbare Graph-Edge übernommen", Strength: .86, Metadata: map[string]any{"trigger": trigger, "relation_type": safeRelation(decision.RelationType), "confidence": decision.Confidence, "semantic_similarity": sim, "research_result_count": len(researchResults), "topic_label": decision.TopicLabel, "candidate_source": candidateSource}})
|
|
if e.RuntimeSettings().ProcessingMode == "clustered" && e.Cfg.ClusterArticleBatching {
|
|
outcome.PendingArticle = &pendingArticleCandidate{Seeds: []model.Node{a, b}, Relation: decision, Research: researchResults}
|
|
e.Broker.Publish(model.Activity{Type: "article.cluster.deferred", Source: "brain", Phase: "knowledge-planning", NodeIDs: []string{a.ID, b.ID}, Message: "Relation wird bis zum Zyklusende mit thematisch ähnlichen Relationen zu einem gemeinsamen Artikelauftrag gebündelt", Strength: .38, Metadata: map[string]any{"trigger": trigger, "topic_label": decision.TopicLabel, "processing_mode": "clustered"}})
|
|
} else {
|
|
articleOutcome, err := e.synthesizeKnowledgeArticle(ctx, trigger, []model.Node{a, b}, decision, researchResults)
|
|
if err != nil {
|
|
// synthesizeKnowledgeArticle owns the terminal article.failed event and
|
|
// its native run_id. Do not publish a second orphan terminal here.
|
|
outcome.ArticleSkipped = true
|
|
} else {
|
|
outcome.ArticleCreated = articleOutcome.Created
|
|
outcome.ArticleSkipped = articleOutcome.Skipped
|
|
}
|
|
}
|
|
} else {
|
|
outcome.Rejected = true
|
|
e.Broker.Publish(model.Activity{Type: "think.rejected", Source: "brain", Phase: "validation", NodeIDs: []string{a.ID, b.ID}, Message: "Ähnlichkeit geprüft, aber nicht als belastbare Edge übernommen", Strength: .42, Metadata: map[string]any{"trigger": trigger, "relation_type": safeRelation(decision.RelationType), "confidence": decision.Confidence, "semantic_similarity": sim, "explanation": decision.Explanation}})
|
|
}
|
|
return outcome, nil
|
|
}
|
|
|
|
func (e *Engine) addResearch(ctx context.Context, a, b model.Node, results []model.ResearchResult) researchGraphRefs {
|
|
refs := researchGraphRefs{}
|
|
for _, r := range results {
|
|
id := graph.ID("external", r.URL)
|
|
metadata := map[string]any{"source": graph.SourceFromURL(r.URL), "query_pair": []string{a.ID, b.ID}, "validation_state": "relation_fulltext_gate_v10", "relevance": r.Relevance, "source_quality": r.SourceQuality, "source_quality_score": r.SourceQualityScore, "assessment_reason": r.AssessmentReason}
|
|
if r.Fetched && r.Relevant {
|
|
if relPath, contentHash, err := e.queueResearchEvidence(r); err == nil {
|
|
metadata["evidence_path"] = relPath
|
|
metadata["evidence_sha256"] = contentHash
|
|
}
|
|
}
|
|
n := model.Node{ID: id, Kind: "external", Label: r.Title, Summary: clamp(r.Content, 700), Status: "research", Origin: "research", ExternalID: r.URL, URI: r.URL, Categories: unique(append(append([]string{}, a.Categories...), b.Categories...)), Weight: math.Max(.8, r.SourceQualityScore), Metadata: metadata, UpdatedAt: time.Now().UTC()}
|
|
refs.Mutations.Add(e.Graph.UpsertNodeWithStats(n))
|
|
refs.NodeIDs = append(refs.NodeIDs, id)
|
|
for _, targetID := range []string{a.ID, b.ID} {
|
|
edge := model.Edge{Source: id, Target: targetID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: math.Max(.65, r.Relevance), Weight: math.Max(.5, r.SourceQualityScore*.7), Explanation: r.AssessmentReason, Metadata: map[string]any{"validation_state": "relation_fulltext_gate_v10", "source_quality": r.SourceQuality, "relevance": r.Relevance}}
|
|
refs.Mutations.Add(e.Graph.UpsertEdgeWithStats(edge))
|
|
refs.EdgeIDs = append(refs.EdgeIDs, graph.EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin))
|
|
}
|
|
}
|
|
refs = uniqueResearchRefs(refs)
|
|
refs.Mutations.Add(e.learnRelationResearchNodes(ctx, refs.NodeIDs))
|
|
return refs
|
|
}
|
|
|
|
// learnRelationResearchNodes embeds accepted relation evidence in the workflow
|
|
// that created it. The scheduled knowledge scanner intentionally does not own
|
|
// external research nodes, otherwise parallel Security/Research work becomes
|
|
// impossible to attribute and external evidence can remain permanently
|
|
// unembedded. A failed embedding is visible through readiness instead of being
|
|
// hidden behind a 256D fallback vector.
|
|
func (e *Engine) learnRelationResearchNodes(ctx context.Context, nodeIDs []string) graph.MutationStats {
|
|
var stats graph.MutationStats
|
|
if !e.LearningEnabled() || e.Ollama == nil || len(nodeIDs) == 0 {
|
|
return stats
|
|
}
|
|
ids := make([]string, 0, len(nodeIDs))
|
|
texts := make([]string, 0, len(nodeIDs))
|
|
for _, id := range unique(nodeIDs) {
|
|
if _, ok := e.Graph.Vector(id); ok {
|
|
continue
|
|
}
|
|
node, ok := e.Graph.GetNode(id)
|
|
if !ok || strings.TrimSpace(embeddingText(node)) == "" {
|
|
continue
|
|
}
|
|
ids = append(ids, id)
|
|
texts = append(texts, embeddingText(node))
|
|
}
|
|
if len(ids) == 0 {
|
|
return stats
|
|
}
|
|
started := time.Now()
|
|
vectors, err := e.Ollama.Embed(ctx, texts)
|
|
if err != nil || len(vectors) != len(ids) {
|
|
e.Broker.Publish(model.Activity{Type: "research.embedding.failed", Source: "ollama", Phase: "embedding", NodeIDs: ids, Message: "Akzeptierte Relationsbelege konnten nicht eingebettet werden; Readiness bleibt bis zum erfolgreichen Retry rot", Strength: .42, Metadata: map[string]any{"result_count": len(ids), "error": errorString(err), "duration_ms": time.Since(started).Milliseconds()}})
|
|
return stats
|
|
}
|
|
for i, vector := range vectors {
|
|
if len(vector) == 0 {
|
|
e.Broker.Publish(model.Activity{Type: "research.embedding.failed", Source: "ollama", Phase: "embedding", NodeIDs: []string{ids[i]}, Message: "Akzeptierter Relationsbeleg erhielt einen leeren Embedding-Vektor", Strength: .42, Metadata: map[string]any{"result_count": 1, "duration_ms": time.Since(started).Milliseconds()}})
|
|
continue
|
|
}
|
|
stats.Add(e.Graph.SetVectorWithStats(ids[i], vector))
|
|
}
|
|
e.Broker.Publish(model.Activity{Type: "research.learned", Source: "ollama", Phase: "embedding", NodeIDs: ids, Message: fmt.Sprintf("%d akzeptierte Relationsbelege wurden unmittelbar eingebettet", int(stats.VectorsCreated+stats.VectorsUpdated)), Strength: .58, Metadata: withRunMutations(map[string]any{"result_count": len(ids), "model": e.Cfg.EmbeddingModel, "duration_ms": time.Since(started).Milliseconds()}, stats)})
|
|
return stats
|
|
}
|
|
func (e *Engine) Status() map[string]any {
|
|
nodes, edges, version := e.Graph.Counts()
|
|
e.stateMu.RLock()
|
|
status := map[string]any{
|
|
"ok": true, "nodes": nodes, "edges": edges, "version": version,
|
|
"last_scan": e.lastScan, "bootstrap_complete": e.bootstrapComplete, "bootstrap_completed_at": e.bootstrapCompletedAt, "bootstrap_error": e.bootstrapError, "last_enrich": e.lastEnrich, "last_enrich_attempt": e.lastAttempt,
|
|
"next_enrich": e.nextEnrich, "ollama_ok": e.ollamaOK, "auto_enrich": e.Cfg.AutoEnrich,
|
|
"enrich_running": e.enrichRunning, "enrich_trigger": e.enrichTrigger, "enrich_result": e.enrichResult,
|
|
"enrich_error": e.enrichError, "enrich_cycles": e.enrichCycles, "enrich_created": e.enrichCreated,
|
|
"enrich_rejected": e.enrichRejected, "relations_created": e.relationsCreated, "articles_created": e.articlesCreated,
|
|
"articles_skipped": e.articlesSkipped, "article_synthesis_enabled": e.Cfg.ArticleSynthesisEnabled,
|
|
"article_cpu_quality_enabled": e.Cfg.ArticleCPUQualityEnabled, "article_cpu_quality_agent_offload": e.Cfg.ArticleCPUQualityAgentOffload, "article_cpu_quality_agent_required": e.Cfg.ArticleCPUQualityAgentRequired, "article_cpu_quality_agent_wait": e.Cfg.ArticleCPUQualityAgentWait.String(),
|
|
"article_min_sources": e.Cfg.ArticleMinSources, "article_max_sources": e.Cfg.ArticleMaxSources,
|
|
"article_min_production_ratio": e.Cfg.ArticleMinProductionRatio, "article_max_generation_depth": e.Cfg.ArticleMaxGenerationDepth,
|
|
"article_max_research_queries": e.Cfg.ArticleMaxResearchQueries, "article_research_results": e.Cfg.ArticleResearchResults,
|
|
"article_research_rounds": e.Cfg.ArticleResearchRounds, "article_research_fetch_results": e.Cfg.ArticleResearchFetchResults,
|
|
"article_research_exploration_results": e.Cfg.ArticleResearchExplorationResults, "article_research_prefetch_min_relevance": e.Cfg.ArticleResearchPrefetchMinRelevance,
|
|
"article_research_min_relevance": e.Cfg.ArticleResearchMinRelevance, "article_research_min_quality": e.Cfg.ArticleResearchMinQuality,
|
|
"article_research_page_max_bytes": e.Cfg.ArticleResearchPageMaxBytes, "article_research_page_max_chars": e.Cfg.ArticleResearchPageMaxChars,
|
|
"article_research_fetch_timeout": e.Cfg.ArticleResearchFetchTimeout.String(), "article_research_allow_private": e.Cfg.ArticleResearchAllowPrivate,
|
|
"article_language": e.Cfg.ArticleLanguage, "article_synthesis_model": e.Cfg.ArticleSynthesisModel, "article_review_model": e.Cfg.ArticleReviewModel, "article_review_repair_rounds": e.Cfg.ArticleReviewRepairRounds, "article_pipeline": "adaptive_generate_review/v5-quality-gate-v12", "article_research_strategy": e.Cfg.ArticleResearchStrategy, "article_effective_research_strategy": e.effectiveArticleResearchStrategy(), "article_adaptive_initial_queries": e.Cfg.ArticleAdaptiveInitialQueries, "article_adaptive_initial_fetch": e.Cfg.ArticleAdaptiveInitialFetch, "research_dedupe": e.researchDedupeStatus(),
|
|
"scan_interval": e.Cfg.ScanInterval.String(), "knowledge_full_verify_interval": e.Cfg.KnowledgeFullVerifyInterval.String(),
|
|
"enrich_interval": e.Cfg.EnrichInterval.String(),
|
|
"enrich_batch_size": e.Cfg.EnrichBatchSize, "enrich_anchors": e.Cfg.EnrichAnchors,
|
|
"processing_mode": e.RuntimeSettings().ProcessingMode, "speed_mode": e.RuntimeSettings().SpeedMode, "speed_cpu_tasks": e.RuntimeSettings().SpeedCPUWorkers, "speed_gpu_tasks": e.RuntimeSettings().SpeedGPUInflight, "speed_algorithm": speedModeVersion, "cluster_hash_bits": e.Cfg.ClusterHashBits, "cluster_hash_tables": e.Cfg.ClusterHashTables,
|
|
"cluster_candidates_per_anchor": e.Cfg.ClusterCandidatesPerAnchor, "cluster_article_candidates": e.Cfg.ClusterArticleCandidates,
|
|
"cluster_review_evidence": e.Cfg.ClusterReviewEvidence, "cluster_review_context_chars": e.Cfg.ClusterReviewContextChars, "cluster_article_batching": e.Cfg.ClusterArticleBatching,
|
|
"vector_graph_enabled": e.Cfg.VectorGraphEnabled, "vector_graph_neighbors": e.Cfg.VectorGraphNeighbors, "vector_graph_candidates": e.Cfg.VectorGraphCandidates, "vector_graph_min_similarity": e.Cfg.VectorGraphMinSimilarity, "vector_graph_min_affinity": e.Cfg.VectorGraphMinAffinity, "vector_graph_layout": e.Cfg.VectorGraphLayout, "vector_graph_layout_mode": e.effectiveVectorLayoutMode(), "vector_graph_reevaluate_interval": e.Cfg.VectorGraphReevaluateInterval.String(), "vector_graph_relax_layout": e.Cfg.VectorGraphRelaxLayout, "vector_graph_layout_relax_effective": e.Cfg.VectorGraphRelaxLayout && !e.Cfg.VectorGraphLayout, "vector_graph_layout_relax_interval": e.Cfg.VectorGraphLayoutRelaxInterval.String(), "vector_graph_layout_blend": e.Cfg.VectorGraphLayoutBlend, "vector_graph_layout_max_shift": e.Cfg.VectorGraphLayoutMaxShift, "last_vector_graph": e.lastVectorGraph, "last_vector_layout": e.lastVectorLayout,
|
|
"vector_graph_orphan_pass": e.Cfg.VectorGraphOrphanPass, "vector_graph_orphan_neighbors": e.Cfg.VectorGraphOrphanNeighbors, "vector_graph_orphan_candidates": e.Cfg.VectorGraphOrphanCandidates, "vector_graph_orphan_min_similarity": e.Cfg.VectorGraphOrphanMinSimilarity, "vector_graph_orphan_min_affinity": e.Cfg.VectorGraphOrphanMinAffinity,
|
|
"vector_graph_agent_offload": e.Cfg.VectorGraphAgentOffload, "vector_graph_agent_required": e.Cfg.VectorGraphAgentRequired, "vector_graph_agent_wait": e.Cfg.VectorGraphAgentWait.String(), "thinking_vector_guided": e.Cfg.ThinkingVectorGuided,
|
|
"research_enabled": e.ResearchEnabledForRuntime(), "chat_model": e.Cfg.ChatModel, "embedding_model": e.Cfg.EmbeddingModel,
|
|
"searxng": e.ResearchStatus(),
|
|
"ollama_pool": e.Ollama.PoolStatus(), "article_model_status": map[string]any{"synthesis": e.Ollama.ModelStatus(e.Cfg.ArticleSynthesisModel), "review": e.Ollama.ModelStatus(e.Cfg.ArticleReviewModel)}, "persistence": e.Persistence.Status(), "graph_storage": e.Graph.StorageStatus(),
|
|
"runtime_settings": e.RuntimeSettingsView(),
|
|
}
|
|
if e.GLPIKB != nil {
|
|
status["glpi_kb"] = e.GLPIKB.Status()
|
|
} else {
|
|
status["glpi_kb"] = ingest.GLPIKBStatus{Enabled: false}
|
|
}
|
|
e.stateMu.RUnlock()
|
|
status["autonomous_research"] = e.AutonomousResearchStatus(context.Background())
|
|
if e.SourceInbox != nil {
|
|
if inbox, err := e.SourceInbox.Stats(context.Background()); err == nil {
|
|
status["source_inbox"] = inbox
|
|
}
|
|
status["source_inbox_classifier"] = map[string]any{
|
|
"version": SourceInboxClassifierVersion,
|
|
"min_similarity": e.Cfg.SourceInboxMinSimilarity,
|
|
"min_priority": e.Cfg.SourceInboxMinPriority,
|
|
"novelty_floor": e.Cfg.SourceInboxNoveltyFloor,
|
|
"security_proactive": e.Cfg.SourceInboxSecurityProactiveEnabled,
|
|
"security_batch_size": e.Cfg.SourceInboxSecurityBatchSize,
|
|
"security_min_priority": e.Cfg.SourceInboxSecurityMinPriority,
|
|
"security_min_confidence": e.Cfg.SourceInboxSecurityMinConfidence,
|
|
"security_fetch_min_chars": e.Cfg.SourceInboxSecurityFetchMinChars,
|
|
"security_research_results": e.Cfg.SourceInboxSecurityResearchResults,
|
|
}
|
|
}
|
|
return status
|
|
}
|
|
|
|
func (e *Engine) Flush(ctx context.Context) error {
|
|
return e.Persistence.Flush(ctx, "manual")
|
|
}
|
|
|
|
func (e *Engine) ExportGraph(ctx context.Context, destination string) error {
|
|
if err := e.Persistence.Flush(ctx, "export"); err != nil {
|
|
return err
|
|
}
|
|
return e.Graph.Export(ctx, destination)
|
|
}
|
|
|
|
func (e *Engine) SyncGLPIKB(ctx context.Context) error {
|
|
if !e.LearningEnabled() {
|
|
return ErrLearningDisabled
|
|
}
|
|
if e.GLPIKB == nil {
|
|
return fmt.Errorf("GLPI knowledge-base integration is disabled")
|
|
}
|
|
return e.GLPIKB.Sync(ctx, "manual")
|
|
}
|
|
|
|
func (e *Engine) relationResearchNeeded(a, b model.Node, similarity float64, decision model.RelationDecision) bool {
|
|
if !decision.NeedsResearch || strings.TrimSpace(decision.ResearchQuery) == "" {
|
|
return false
|
|
}
|
|
relationType := safeRelation(decision.RelationType)
|
|
// same_topic/related_to are graph-topology judgements. For two internal
|
|
// knowledge entries with very high semantic overlap and a confident model
|
|
// decision, external web evidence does not make the relationship more true;
|
|
// it only adds cost and unrelated research nodes. Keep web verification for
|
|
// factual/dependency/contradiction relations and for freshness-sensitive
|
|
// language.
|
|
if decision.Related && decision.Confidence >= e.Cfg.RelationThreshold && similarity >= .90 && (relationType == "same_topic" || relationType == "related_to") {
|
|
context := strings.Join([]string{decision.ResearchQuery, decision.Explanation, a.Label, b.Label}, " ")
|
|
if !containsFreshnessLanguage(context) {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
func relationContextWithResearch(a, b model.Node, sim float64, results []model.ResearchResult) string {
|
|
var out strings.Builder
|
|
out.WriteString(relationContext(a, b, sim))
|
|
out.WriteString("\n\nWEB-SUCHERGEBNISSE (ungeprüfte Hinweise):\n")
|
|
for i, r := range results {
|
|
fmt.Fprintf(&out, "\n%d. %s\nURL: %s\nAuszug: %s\n", i+1, r.Title, r.URL, clamp(r.Content, 700))
|
|
}
|
|
return out.String()
|
|
}
|
|
|
|
func relationContext(a, b model.Node, sim float64) string {
|
|
// Internal IDs are intentionally omitted from the model-visible relation
|
|
// context. They are routing metadata, not semantic search terms, and older
|
|
// prompts occasionally copied them into public SearXNG queries.
|
|
return fmt.Sprintf("SEMANTISCHE_NÄHE: %.4f\n\nA\nTitel: %s\nKategorien: %s\nInhalt: %s\n\nB\nTitel: %s\nKategorien: %s\nInhalt: %s", sim, a.Label, strings.Join(a.Categories, ", "), a.Summary, b.Label, strings.Join(b.Categories, ", "), b.Summary)
|
|
}
|
|
func relationSchema() map[string]any {
|
|
return map[string]any{"type": "object", "properties": map[string]any{"related": map[string]any{"type": "boolean"}, "relation_type": map[string]any{"type": "string", "enum": []string{"related_to", "depends_on", "supports", "contradicts", "extends", "same_topic", "caused_by"}}, "confidence": map[string]any{"type": "number", "minimum": 0, "maximum": 1}, "explanation": map[string]any{"type": "string"}, "needs_research": map[string]any{"type": "boolean"}, "research_query": map[string]any{"type": "string"}, "topic_label": map[string]any{"type": "string"}, "keywords": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}}, "required": []string{"related", "relation_type", "confidence", "explanation", "needs_research", "research_query", "topic_label", "keywords"}}
|
|
}
|
|
func answerSchema() map[string]any {
|
|
return map[string]any{"type": "object", "properties": map[string]any{"answer": map[string]any{"type": "string"}, "used_node_ids": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "uncertainties": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}}, "required": []string{"answer", "used_node_ids", "uncertainties"}}
|
|
}
|
|
func validIDs(in, allowed []string) []string {
|
|
set := map[string]bool{}
|
|
for _, x := range allowed {
|
|
set[x] = true
|
|
}
|
|
var out []string
|
|
for _, x := range in {
|
|
if set[x] {
|
|
out = append(out, x)
|
|
}
|
|
}
|
|
if len(out) == 0 {
|
|
return allowed
|
|
}
|
|
return unique(out)
|
|
}
|
|
func safeRelation(s string) string {
|
|
switch s {
|
|
case "related_to", "depends_on", "supports", "contradicts", "extends", "same_topic", "caused_by":
|
|
return s
|
|
default:
|
|
return "related_to"
|
|
}
|
|
}
|
|
func common(a, b []string) []string {
|
|
set := map[string]string{}
|
|
for _, x := range a {
|
|
set[strings.ToLower(x)] = x
|
|
}
|
|
var out []string
|
|
for _, x := range b {
|
|
if v, ok := set[strings.ToLower(x)]; ok {
|
|
out = append(out, v)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
func first(in []string, n int) []string {
|
|
if len(in) > n {
|
|
return in[:n]
|
|
}
|
|
return in
|
|
}
|
|
func unique(in []string) []string {
|
|
set := map[string]bool{}
|
|
var out []string
|
|
for _, x := range in {
|
|
x = strings.TrimSpace(x)
|
|
k := strings.ToLower(x)
|
|
if x == "" || set[k] {
|
|
continue
|
|
}
|
|
set[k] = true
|
|
out = append(out, x)
|
|
}
|
|
sort.Strings(out)
|
|
return out
|
|
}
|
|
func clamp(s string, n int) string {
|
|
r := []rune(strings.TrimSpace(s))
|
|
if len(r) <= n {
|
|
return string(r)
|
|
}
|
|
return string(r[:n]) + "…"
|
|
}
|
|
func nonempty(a, b string) string {
|
|
if strings.TrimSpace(a) != "" {
|
|
return strings.TrimSpace(a)
|
|
}
|
|
return b
|
|
}
|