package sourceagent import ( "context" "math/bits" "net/http" "net/http/httptest" "testing" "time" "github.com/local/glpi-neural-brain/internal/research" ) func TestValidateTaskSupportsPollerTypes(t *testing.T) { for _, typ := range []string{"rss", "atom", "sitemap", "web"} { task, err := validateTask(Task{ID: "security", AgentID: "a", Name: "Security", Type: typ, URL: "https://example.org/feed", Enabled: true, PollInterval: "2h", MaxItems: 25}) if err != nil { t.Fatalf("%s: %v", typ, err) } if task.Type != typ || task.MaxItems != 25 { t.Fatalf("unexpected task: %+v", task) } } } func TestValidateTaskRejectsUnsafeShape(t *testing.T) { if _, err := validateTask(Task{Type: "ftp", URL: "https://example.org", PollInterval: "1h"}); err == nil { t.Fatal("expected unsupported type") } if _, err := validateTask(Task{Type: "rss", URL: "ftp://example.org/feed", PollInterval: "1h"}); err == nil { t.Fatal("expected URL validation error") } if _, err := validateTask(Task{Type: "rss", URL: "https://example.org/feed", PollInterval: "1m"}); err == nil { t.Fatal("expected interval validation error") } if _, err := validateTask(Task{Type: "rss", URL: "https://user:pass@example.org/feed", PollInterval: "1h"}); err == nil { t.Fatal("expected userinfo URL validation error") } } func TestNormalizeDocumentBuildsContentHash(t *testing.T) { d, err := normalizeDocument(Document{URL: "https://example.org/a", Title: "A", Text: "Dies ist ein ausreichend langer Dokumenttext mit sicherheitsrelevanten Informationen und mehreren Details für die Verarbeitung."}) if err != nil { t.Fatal(err) } if d.ContentSHA256 == "" || d.CanonicalURL != "https://example.org/a" || d.ExternalID != "https://example.org/a" { t.Fatalf("unexpected normalized document: %+v", d) } } func TestSimhashKeepsRelatedTextCloser(t *testing.T) { base := simhash64("Windows Backup Repository Ransomware Schutz immutable storage MFA") related := simhash64("Ransomware Schutz für Windows Backup Repository mit MFA und immutable Storage") unrelated := simhash64("Kaffee Bohnen Espresso Maschine Mahlgrad Temperatur") if bits.OnesCount64(base^related) >= bits.OnesCount64(base^unrelated) { t.Fatalf("related text should hash closer") } } func TestFetchRawBlocksPrivateSourceByDefault(t *testing.T) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write([]byte("")) })) defer srv.Close() r := &Runner{cfg: RunnerConfig{HTTPTimeout: time.Second}, sourceHTTP: research.NewSafeHTTPClient(false, time.Second)} if _, _, _, _, err := r.fetchRaw(context.Background(), srv.URL, 1024); err == nil { t.Fatal("expected private/loopback source URL to be blocked") } } func TestSitemapIndexCollectsChildURLs(t *testing.T) { mux := http.NewServeMux() var base string mux.HandleFunc("/index.xml", func(w http.ResponseWriter, _ *http.Request) { w.Header().Set("Content-Type", "application/xml") _, _ = w.Write([]byte(`` + base + `/child.xml`)) }) mux.HandleFunc("/child.xml", func(w http.ResponseWriter, _ *http.Request) { w.Header().Set("Content-Type", "application/xml") _, _ = w.Write([]byte(`` + base + `/article-12026-08-07T12:00:00Z`)) }) srv := httptest.NewServer(mux) defer srv.Close() base = srv.URL r := &Runner{cfg: RunnerConfig{HTTPTimeout: time.Second, AllowPrivate: true}, sourceHTTP: research.NewSafeHTTPClient(true, time.Second)} items, err := r.collectSitemapItems(context.Background(), srv.URL+"/index.xml", 10, 0) if err != nil { t.Fatal(err) } if len(items) != 1 || items[0].Link != srv.URL+"/article-1" { t.Fatalf("unexpected sitemap items: %+v", items) } } func TestRunnerStatusDoesNotTreatCachedConfigAsLiveConnection(t *testing.T) { r := &Runner{ cfg: RunnerConfig{BrainURL: "http://127.0.0.1:8091", AgentID: "agent-test", Version: "test"}, remote: RemoteConfig{IssuedAt: time.Now().UTC(), Tasks: []Task{{ID: "cached-task"}}}, diag: runnerDiagnostics{ConfigSource: "cache", CacheLoadedAt: time.Now().UTC()}, } st := r.Status() if st["brain_connected"] != false { t.Fatalf("cached config must not be reported as live connection: %+v", st) } if st["connection_state"] != "cached_config_only" { t.Fatalf("unexpected state: %+v", st) } if st["brain_url_warning"] == "" { t.Fatalf("expected loopback warning: %+v", st) } } func TestRunnerStatusReportsFreshBrainConnection(t *testing.T) { now := time.Now().UTC() r := &Runner{ cfg: RunnerConfig{BrainURL: "http://brain:8090", AgentID: "agent-test", Version: "test"}, remote: RemoteConfig{IssuedAt: now}, diag: runnerDiagnostics{ConfigSource: "brain", LastConfigAttemptAt: now.Add(-time.Second), LastConfigSuccessAt: now}, } st := r.Status() if st["brain_connected"] != true || st["connection_state"] != "connected" { t.Fatalf("expected live connection: %+v", st) } } func TestRunnerStatusAcceptsConnectedLoopbackForNativeTest(t *testing.T) { now := time.Now().UTC() r := &Runner{ cfg: RunnerConfig{BrainURL: "http://127.0.0.1:8091", AgentID: "agent-test", Version: "test"}, remote: RemoteConfig{IssuedAt: now}, diag: runnerDiagnostics{ConfigSource: "brain", LastConfigSuccessAt: now, LastHeartbeatSuccess: now}, } st := r.Status() if st["brain_connected"] != true { t.Fatalf("expected connected loopback to be accepted: %+v", st) } if st["brain_url_warning"] != "" { t.Fatalf("connected native loopback must not be warned as broken: %+v", st) } } func TestEnsureInboxClassifierVersionRequeuesArchivedOnce(t *testing.T) { ctx := context.Background() store, err := OpenStore(t.TempDir()) if err != nil { t.Fatal(err) } defer store.Close() agent, _, err := store.CreateAgent(ctx, "agent-a", "Agent A") if err != nil { t.Fatal(err) } task, err := store.UpsertTask(ctx, Task{ID: "task-a", AgentID: agent.ID, Name: "Security", Type: "rss", URL: "https://example.org/feed", Enabled: true, PollInterval: "1h", MaxItems: 10}) if err != nil { t.Fatal(err) } _, err = store.Ingest(ctx, agent.ID, task.ID, []Document{{URL: "https://example.org/a", Title: "Security update", Text: "Dies ist ein ausreichend langer Sicherheitsartikel mit mehreren Details, damit das Dokument in die Source Inbox aufgenommen und klassifiziert werden kann."}}) if err != nil { t.Fatal(err) } claimed, err := store.ClaimInbox(ctx, 1) if err != nil || len(claimed) != 1 { t.Fatalf("claim failed: %#v %v", claimed, err) } if err := store.CompleteClassification(ctx, claimed[0].ID, "archived", .4, "kb-1", map[string]any{}); err != nil { t.Fatal(err) } n, err := store.EnsureInboxClassifierVersion(ctx, 2) if err != nil || n != 1 { t.Fatalf("expected one requeued archived item, n=%d err=%v", n, err) } items, err := store.ListInbox(ctx, "received", 10) if err != nil || len(items) != 1 { t.Fatalf("expected requeued received item: %#v %v", items, err) } n, err = store.EnsureInboxClassifierVersion(ctx, 2) if err != nil || n != 0 { t.Fatalf("classifier migration must be one-shot, n=%d err=%v", n, err) } } func TestProactiveSecurityInboxLifecycle(t *testing.T) { ctx := context.Background() store, err := OpenStore(t.TempDir()) if err != nil { t.Fatal(err) } defer store.Close() agent, _, err := store.CreateAgent(ctx, "agent-security", "Security Agent") if err != nil { t.Fatal(err) } task, err := store.UpsertTask(ctx, Task{ID: "task-security", AgentID: agent.ID, Name: "Security Alerts", Type: "rss", URL: "https://example.org/feed", Enabled: true, PollInterval: "1h", MaxItems: 10}) if err != nil { t.Fatal(err) } _, err = store.Ingest(ctx, agent.ID, task.ID, []Document{{URL: "https://example.org/advisory", CanonicalURL: "https://example.org/advisory", Title: "Critical security advisory", Text: "A concrete security advisory with enough source text to be classified and later materialized into a proactive security node.", Metadata: map[string]any{"security_proactive": "true"}}}) if err != nil { t.Fatal(err) } claimed, err := store.ClaimInbox(ctx, 1) if err != nil || len(claimed) != 1 { t.Fatalf("claim classification: %#v %v", claimed, err) } if err := store.CompleteClassification(ctx, claimed[0].ID, "candidate", .71, "kb-security", map[string]any{"priority_score": .71, "security_candidate": true}); err != nil { t.Fatal(err) } if err := store.QueueProactiveSecurity(ctx, claimed[0].ID); err != nil { t.Fatal(err) } securityItems, err := store.ClaimProactiveSecurity(ctx, 1) if err != nil || len(securityItems) != 1 || securityItems[0].ProactiveState != "processing" { t.Fatalf("claim proactive: %#v %v", securityItems, err) } if err := store.CompleteProactiveSecurity(ctx, claimed[0].ID, "external-node-1", map[string]any{"security_event_type": "vulnerability"}); err != nil { t.Fatal(err) } materialized, err := store.ListInbox(ctx, "materialized", 10) if err != nil || len(materialized) != 1 || materialized[0].MaterializedNodeID != "external-node-1" || materialized[0].ProactiveState != "done" { t.Fatalf("materialized lifecycle mismatch: %#v %v", materialized, err) } results, err := store.SearchCandidates(ctx, "critical security advisory", 5, 0) if err != nil || len(results) != 1 { t.Fatalf("materialized document must remain searchable before SearXNG: %#v %v", results, err) } } func TestEffectiveConcurrencyUsesRemoteSpeedLimit(t *testing.T) { r := &Runner{ cfg: RunnerConfig{Concurrency: 3}, remote: RemoteConfig{Performance: AgentPerformance{SpeedMode: true, CPUWorkers: 12}}, } if got := r.effectiveConcurrency(); got != 12 { t.Fatalf("speed concurrency=%d want 12", got) } r.remote.Performance.SpeedMode = false if got := r.effectiveConcurrency(); got != 3 { t.Fatalf("normal concurrency=%d want 3", got) } }