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)
}
}