All checks were successful
release-tag / release-image (push) Successful in 2m32s
239 lines
9.6 KiB
Go
239 lines
9.6 KiB
Go
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("<rss/>")) }))
|
|
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(`<sitemapindex><sitemap><loc>` + base + `/child.xml</loc></sitemap></sitemapindex>`))
|
|
})
|
|
mux.HandleFunc("/child.xml", func(w http.ResponseWriter, _ *http.Request) {
|
|
w.Header().Set("Content-Type", "application/xml")
|
|
_, _ = w.Write([]byte(`<urlset><url><loc>` + base + `/article-1</loc><lastmod>2026-08-07T12:00:00Z</lastmod></url></urlset>`))
|
|
})
|
|
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)
|
|
}
|
|
}
|