Agent-Mode mit Aufgabenteilung und automatischer Recherche

This commit is contained in:
2026-08-07 21:57:23 +02:00
parent ffee925163
commit adbe69d24d
30 changed files with 2580 additions and 51 deletions
+658
View File
@@ -0,0 +1,658 @@
package sourceagent
import (
"bytes"
"context"
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"encoding/xml"
"errors"
"fmt"
"html"
"io"
"log/slog"
"net/http"
"net/url"
"os"
"path/filepath"
"regexp"
"sort"
"strings"
"sync"
"time"
"github.com/local/glpi-neural-brain/internal/research"
_ "modernc.org/sqlite"
)
type RunnerConfig struct {
BrainURL string
AgentID string
Token string
DataDir string
ConfigFile string
ConfigRefresh time.Duration
HTTPTimeout time.Duration
Concurrency int
BatchSize int
AllowPrivate bool
Version string
}
type Runner struct {
cfg RunnerConfig
http *http.Client
sourceHTTP *http.Client
state *localState
mu sync.RWMutex
remote RemoteConfig
wake chan struct{}
}
type bootstrapConfig struct {
BrainURL string `json:"brain_url"`
AgentID string `json:"agent_id"`
Token string `json:"token"`
}
func NewRunner(cfg RunnerConfig) (*Runner, error) {
if strings.TrimSpace(cfg.ConfigFile) != "" {
if data, err := os.ReadFile(cfg.ConfigFile); err == nil {
var b bootstrapConfig
if json.Unmarshal(data, &b) == nil {
if cfg.BrainURL == "" {
cfg.BrainURL = b.BrainURL
}
if cfg.AgentID == "" {
cfg.AgentID = b.AgentID
}
if cfg.Token == "" {
cfg.Token = b.Token
}
}
}
}
cfg.BrainURL = strings.TrimRight(strings.TrimSpace(cfg.BrainURL), "/")
cfg.AgentID = strings.TrimSpace(cfg.AgentID)
cfg.Token = strings.TrimSpace(cfg.Token)
if cfg.BrainURL == "" || cfg.AgentID == "" || cfg.Token == "" {
return nil, errors.New("agent mode requires BRAIN_AGENT_BRAIN_URL, BRAIN_AGENT_ID and BRAIN_AGENT_TOKEN (or BRAIN_AGENT_CONFIG_FILE)")
}
u, err := url.Parse(cfg.BrainURL)
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") {
return nil, errors.New("BRAIN_AGENT_BRAIN_URL must be an absolute http(s) URL")
}
if u.User != nil {
return nil, errors.New("BRAIN_AGENT_BRAIN_URL must not contain userinfo")
}
if cfg.ConfigRefresh < time.Minute {
cfg.ConfigRefresh = 5 * time.Minute
}
if cfg.HTTPTimeout < 5*time.Second {
cfg.HTTPTimeout = 30 * time.Second
}
if cfg.Concurrency < 1 {
cfg.Concurrency = 3
}
if cfg.Concurrency > 16 {
cfg.Concurrency = 16
}
if cfg.BatchSize < 1 {
cfg.BatchSize = 50
}
if cfg.BatchSize > 500 {
cfg.BatchSize = 500
}
state, err := openLocalState(cfg.DataDir)
if err != nil {
return nil, err
}
return &Runner{
cfg: cfg,
http: &http.Client{Timeout: cfg.HTTPTimeout},
sourceHTTP: research.NewSafeHTTPClient(cfg.AllowPrivate, cfg.HTTPTimeout),
state: state,
wake: make(chan struct{}, 1),
}, nil
}
func (r *Runner) Close() error {
if r == nil || r.state == nil {
return nil
}
return r.state.Close()
}
func (r *Runner) Start(ctx context.Context) {
r.loadCachedConfig()
go r.loop(ctx)
}
func (r *Runner) Status() map[string]any {
r.mu.RLock()
remote := r.remote
r.mu.RUnlock()
return map[string]any{"ok": true, "mode": "agent", "agent_id": r.cfg.AgentID, "brain_url": r.cfg.BrainURL, "configured_tasks": len(remote.Tasks), "config_issued_at": remote.IssuedAt, "version": r.cfg.Version}
}
func (r *Runner) loop(ctx context.Context) {
refresh := time.NewTicker(r.cfg.ConfigRefresh)
defer refresh.Stop()
run := time.NewTicker(30 * time.Second)
defer run.Stop()
_ = r.refreshConfig(ctx)
r.runDue(ctx)
for {
select {
case <-ctx.Done():
return
case <-refresh.C:
_ = r.refreshConfig(ctx)
case <-run.C:
r.runDue(ctx)
case <-r.wake:
r.runDue(ctx)
}
}
}
func (r *Runner) refreshConfig(ctx context.Context) error {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/config", nil)
if err != nil {
return err
}
r.auth(req)
resp, err := r.http.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return fmt.Errorf("brain config returned HTTP %d", resp.StatusCode)
}
var cfg RemoteConfig
if err := json.NewDecoder(io.LimitReader(resp.Body, 4<<20)).Decode(&cfg); err != nil {
return err
}
if cfg.Agent.ID != "" && cfg.Agent.ID != r.cfg.AgentID {
return fmt.Errorf("brain returned config for unexpected agent %q", cfg.Agent.ID)
}
r.mu.Lock()
r.remote = cfg
r.mu.Unlock()
r.saveCachedConfig(cfg)
_ = r.sendHeartbeat(ctx, Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "online", Metadata: map[string]any{"configured_tasks": len(cfg.Tasks)}})
return nil
}
func (r *Runner) runDue(ctx context.Context) {
r.mu.RLock()
tasks := append([]Task(nil), r.remote.Tasks...)
r.mu.RUnlock()
if len(tasks) == 0 {
return
}
sem := make(chan struct{}, r.cfg.Concurrency)
var wg sync.WaitGroup
for _, task := range tasks {
task := task
if !task.Enabled || !r.state.Due(ctx, task) {
continue
}
wg.Add(1)
go func() {
defer wg.Done()
select {
case sem <- struct{}{}:
case <-ctx.Done():
return
}
defer func() { <-sem }()
r.runTask(ctx, task)
}()
}
wg.Wait()
}
func (r *Runner) runTask(ctx context.Context, task Task) {
started := time.Now().UTC()
docs, err := r.pollTask(ctx, task)
if err == nil && len(docs) > 0 {
for start := 0; start < len(docs); start += r.cfg.BatchSize {
end := start + r.cfg.BatchSize
if end > len(docs) {
end = len(docs)
}
if sendErr := r.sendBatch(ctx, task.ID, docs[start:end]); sendErr != nil {
err = sendErr
break
}
for _, d := range docs[start:end] {
_ = r.state.MarkSeen(ctx, task.ID, d.CanonicalURL, d.ContentSHA256)
}
}
}
_ = r.state.FinishTask(ctx, task.ID, started, err)
h := Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "ok", LastRunAt: started, TasksChecked: 1, Documents: len(docs)}
if err != nil {
h.Status = "error"
h.LastError = err.Error()
slog.Warn("source agent task failed", "task", task.ID, "error", err)
} else {
slog.Info("source agent task completed", "task", task.ID, "documents", len(docs))
}
_ = r.sendHeartbeat(ctx, h)
}
func (r *Runner) pollTask(ctx context.Context, task Task) ([]Document, error) {
switch task.Type {
case "rss", "atom":
return r.pollFeed(ctx, task)
case "sitemap":
return r.pollSitemap(ctx, task)
case "web":
return r.pollWeb(ctx, task)
default:
return nil, fmt.Errorf("unsupported task type %q", task.Type)
}
}
type feedEnvelope struct {
Channel struct {
Items []struct {
Title string `xml:"title"`
Link string `xml:"link"`
Description string `xml:"description"`
PubDate string `xml:"pubDate"`
GUID string `xml:"guid"`
} `xml:"item"`
} `xml:"channel"`
Entries []struct {
Title string `xml:"title"`
ID string `xml:"id"`
Updated string `xml:"updated"`
Published string `xml:"published"`
Summary string `xml:"summary"`
Content string `xml:"content"`
Links []struct {
Href string `xml:"href,attr"`
Rel string `xml:"rel,attr"`
} `xml:"link"`
} `xml:"entry"`
}
type feedItem struct {
Title, Link, Summary string
Published time.Time
}
func (r *Runner) pollFeed(ctx context.Context, task Task) ([]Document, error) {
body, finalURL, _, _, err := r.fetchRaw(ctx, task.URL, 5<<20)
if err != nil {
return nil, err
}
var env feedEnvelope
if err := xml.Unmarshal(body, &env); err != nil {
return nil, fmt.Errorf("feed XML: %w", err)
}
var items []feedItem
for _, it := range env.Channel.Items {
link := strings.TrimSpace(it.Link)
if link == "" {
link = strings.TrimSpace(it.GUID)
}
items = append(items, feedItem{Title: strings.TrimSpace(it.Title), Link: resolveLink(finalURL, link), Summary: stripMarkup(it.Description), Published: parsePublished(it.PubDate)})
}
for _, it := range env.Entries {
link := ""
for _, l := range it.Links {
if l.Rel == "" || l.Rel == "alternate" {
link = l.Href
break
}
}
pub := parsePublished(it.Published)
if pub.IsZero() {
pub = parsePublished(it.Updated)
}
summary := it.Summary
if summary == "" {
summary = it.Content
}
items = append(items, feedItem{Title: strings.TrimSpace(it.Title), Link: resolveLink(finalURL, link), Summary: stripMarkup(summary), Published: pub})
}
sort.SliceStable(items, func(i, j int) bool { return items[i].Published.After(items[j].Published) })
if len(items) > task.MaxItems {
items = items[:task.MaxItems]
}
return r.materializeItems(ctx, task, items)
}
type sitemapEnvelope struct {
URLs []struct {
Loc string `xml:"loc"`
LastMod string `xml:"lastmod"`
} `xml:"url"`
Sitemaps []struct {
Loc string `xml:"loc"`
} `xml:"sitemap"`
}
func (r *Runner) pollSitemap(ctx context.Context, task Task) ([]Document, error) {
items, err := r.collectSitemapItems(ctx, task.URL, task.MaxItems, 0)
if err != nil {
return nil, err
}
return r.materializeItems(ctx, task, items)
}
func (r *Runner) collectSitemapItems(ctx context.Context, rawURL string, limit, depth int) ([]feedItem, error) {
if limit <= 0 || depth > 1 {
return nil, nil
}
body, finalURL, _, _, err := r.fetchRaw(ctx, rawURL, 8<<20)
if err != nil {
return nil, err
}
var sm sitemapEnvelope
if err := xml.Unmarshal(body, &sm); err != nil {
return nil, err
}
items := make([]feedItem, 0, limit)
for _, u := range sm.URLs {
link := resolveLink(finalURL, u.Loc)
if link == "" {
continue
}
items = append(items, feedItem{Link: link, Published: parsePublished(u.LastMod)})
if len(items) >= limit {
return items, nil
}
}
// Sitemap indexes are common on larger publishers. Follow a bounded number of
// child maps once; source HTTP safety rules apply to every child URL.
for i, child := range sm.Sitemaps {
if len(items) >= limit || i >= 12 {
break
}
childURL := resolveLink(finalURL, child.Loc)
if childURL == "" {
continue
}
more, childErr := r.collectSitemapItems(ctx, childURL, limit-len(items), depth+1)
if childErr != nil {
continue
}
items = append(items, more...)
}
return items, nil
}
var hrefPattern = regexp.MustCompile(`(?is)<a\b[^>]*href\s*=\s*["']([^"'#]+)["'][^>]*>(.*?)</a>`)
func (r *Runner) pollWeb(ctx context.Context, task Task) ([]Document, error) {
body, finalURL, _, _, err := r.fetchRaw(ctx, task.URL, 5<<20)
if err != nil {
return nil, err
}
matches := hrefPattern.FindAllStringSubmatch(string(body), -1)
seen := map[string]bool{}
items := make([]feedItem, 0, task.MaxItems)
base, _ := url.Parse(finalURL)
for _, m := range matches {
link := resolveLink(finalURL, m[1])
if link == "" || seen[link] {
continue
}
u, err := url.Parse(link)
if err != nil || u.Host != base.Host {
continue
}
if !looksArticleLink(u.Path, stripMarkup(m[2])) {
continue
}
seen[link] = true
items = append(items, feedItem{Title: stripMarkup(m[2]), Link: link})
if len(items) >= task.MaxItems {
break
}
}
return r.materializeItems(ctx, task, items)
}
func (r *Runner) materializeItems(ctx context.Context, task Task, items []feedItem) ([]Document, error) {
fetcher := research.New("")
out := make([]Document, 0, len(items))
for _, item := range items {
if strings.TrimSpace(item.Link) == "" {
continue
}
if r.state.SeenURL(ctx, task.ID, item.Link) && !strings.EqualFold(strings.TrimSpace(task.Config["refetch_seen"]), "true") {
continue
}
page, diag, err := fetcher.FetchPage(ctx, item.Link, research.FetchOptions{MaxBytes: 2 << 20, MaxChars: 20000, Timeout: r.cfg.HTTPTimeout, AllowPrivate: r.cfg.AllowPrivate})
text := strings.TrimSpace(item.Summary)
title := strings.TrimSpace(item.Title)
ctype := "text/html"
final := item.Link
if err == nil {
if page.Content != "" {
text = page.Content
}
if page.Title != "" {
title = page.Title
}
if page.URL != "" {
final = page.URL
}
ctype = page.ContentType
} else if len([]rune(text)) < 80 {
continue
}
if title == "" {
title = final
}
sum := sha256.Sum256([]byte(text))
sha := hex.EncodeToString(sum[:])
if r.state.SeenHash(ctx, task.ID, final, sha) {
continue
}
baseURL := ""
if u, e := url.Parse(task.URL); e == nil {
baseURL = u.Scheme + "://" + u.Host
}
out = append(out, Document{ExternalID: final, URL: final, CanonicalURL: final, Title: title, PublishedAt: item.Published, DiscoveredAt: time.Now().UTC(), ContentType: ctype, Text: text, ContentSHA256: sha, SourceName: task.Name, SourceBaseURL: baseURL, Categories: task.Categories, Metadata: map[string]any{"agent_task_type": task.Type, "fetch_error_kind": diag.ErrorKind}})
if len(out) >= task.MaxItems {
break
}
}
return out, nil
}
func (r *Runner) fetchRaw(ctx context.Context, raw string, max int64) ([]byte, string, string, string, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, raw, nil)
if err != nil {
return nil, "", "", "", err
}
req.Header.Set("User-Agent", "glpi-neural-brain-source-agent/1.0")
req.Header.Set("Accept", "application/rss+xml, application/atom+xml, application/xml, text/xml, text/html;q=0.9, */*;q=0.5")
resp, err := r.sourceHTTP.Do(req)
if err != nil {
return nil, "", "", "", err
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return nil, "", "", "", fmt.Errorf("source returned HTTP %d", resp.StatusCode)
}
body, err := io.ReadAll(io.LimitReader(resp.Body, max+1))
if err != nil {
return nil, "", "", "", err
}
if int64(len(body)) > max {
return nil, "", "", "", errors.New("source response too large")
}
return body, resp.Request.URL.String(), resp.Header.Get("ETag"), resp.Header.Get("Last-Modified"), nil
}
func (r *Runner) sendBatch(ctx context.Context, taskID string, docs []Document) error {
payload := IngestBatch{SchemaVersion: SchemaVersion, AgentID: r.cfg.AgentID, TaskID: taskID, Documents: docs}
data, _ := json.Marshal(payload)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.cfg.BrainURL+"/api/v1/agent/ingest", bytes.NewReader(data))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
r.auth(req)
resp, err := r.http.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
return fmt.Errorf("brain ingest HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
}
return nil
}
func (r *Runner) sendHeartbeat(ctx context.Context, h Heartbeat) error {
data, _ := json.Marshal(h)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.cfg.BrainURL+"/api/v1/agent/heartbeat", bytes.NewReader(data))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
r.auth(req)
resp, err := r.http.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return fmt.Errorf("heartbeat HTTP %d", resp.StatusCode)
}
return nil
}
func (r *Runner) auth(req *http.Request) {
req.Header.Set("Authorization", "Bearer "+r.cfg.Token)
req.Header.Set("X-Brain-Agent-ID", r.cfg.AgentID)
}
func resolveLink(base, ref string) string {
ref = strings.TrimSpace(html.UnescapeString(ref))
if ref == "" {
return ""
}
u, err := url.Parse(ref)
if err != nil {
return ""
}
b, err := url.Parse(base)
if err != nil {
return ""
}
return b.ResolveReference(u).String()
}
func stripMarkup(v string) string {
v = regexp.MustCompile(`(?is)<[^>]+>`).ReplaceAllString(v, " ")
return strings.Join(strings.Fields(html.UnescapeString(v)), " ")
}
func looksArticleLink(path, label string) bool {
p := strings.ToLower(path + " " + label)
if strings.Contains(p, "/tag/") || strings.Contains(p, "/category/") || strings.Contains(p, "/author/") || strings.Contains(p, "login") || strings.Contains(p, "privacy") || strings.Contains(p, "impress") || strings.Contains(p, "kontakt") {
return false
}
segments := strings.Split(strings.Trim(path, "/"), "/")
return len(segments) >= 2 || strings.Contains(p, "news") || strings.Contains(p, "blog") || strings.Contains(p, "advis") || strings.Contains(p, "release") || strings.Contains(p, "security")
}
func parsePublished(v string) time.Time {
v = strings.TrimSpace(v)
for _, layout := range []string{time.RFC3339, time.RFC1123Z, time.RFC1123, time.RFC822Z, time.RFC822, "2006-01-02"} {
if t, err := time.Parse(layout, v); err == nil {
return t.UTC()
}
}
return time.Time{}
}
func (r *Runner) configCachePath() string {
return filepath.Join(r.cfg.DataDir, "source-agent-config-cache.json")
}
func (r *Runner) saveCachedConfig(cfg RemoteConfig) {
data, err := json.MarshalIndent(cfg, "", " ")
if err != nil {
return
}
_ = os.WriteFile(r.configCachePath(), append(data, '\n'), 0o600)
}
func (r *Runner) loadCachedConfig() {
data, err := os.ReadFile(r.configCachePath())
if err != nil {
return
}
var cfg RemoteConfig
if json.Unmarshal(data, &cfg) != nil || cfg.Agent.ID != r.cfg.AgentID {
return
}
r.mu.Lock()
r.remote = cfg
r.mu.Unlock()
}
type localState struct {
db *sql.DB
mu sync.Mutex
}
func openLocalState(dataDir string) (*localState, error) {
if strings.TrimSpace(dataDir) == "" {
dataDir = "./data"
}
if err := os.MkdirAll(dataDir, 0o750); err != nil {
return nil, err
}
db, err := sql.Open("sqlite", "file:"+filepath.ToSlash(filepath.Join(dataDir, "source-agent-local.db"))+"?_pragma=busy_timeout(5000)&_pragma=journal_mode(WAL)")
if err != nil {
return nil, err
}
st := &localState{db: db}
for _, q := range []string{`CREATE TABLE IF NOT EXISTS task_state(task_id TEXT PRIMARY KEY,last_run_ns INTEGER NOT NULL DEFAULT 0,last_error TEXT NOT NULL DEFAULT '') WITHOUT ROWID`, `CREATE TABLE IF NOT EXISTS seen(task_id TEXT NOT NULL,url TEXT NOT NULL,content_sha256 TEXT NOT NULL,seen_at_ns INTEGER NOT NULL,PRIMARY KEY(task_id,url,content_sha256)) WITHOUT ROWID`, `CREATE INDEX IF NOT EXISTS idx_seen_task_url ON seen(task_id,url,seen_at_ns DESC)`} {
if _, err := db.Exec(q); err != nil {
db.Close()
return nil, err
}
}
return st, nil
}
func (s *localState) Close() error { return s.db.Close() }
func (s *localState) Due(ctx context.Context, t Task) bool {
d, err := time.ParseDuration(t.PollInterval)
if err != nil {
d = 4 * time.Hour
}
var last int64
err = s.db.QueryRowContext(ctx, `SELECT last_run_ns FROM task_state WHERE task_id=?`, t.ID).Scan(&last)
return err == sql.ErrNoRows || err != nil || last == 0 || time.Since(time.Unix(0, last)) >= d
}
func (s *localState) FinishTask(ctx context.Context, id string, started time.Time, err error) error {
msg := ""
if err != nil {
msg = err.Error()
}
_, e := s.db.ExecContext(ctx, `INSERT INTO task_state(task_id,last_run_ns,last_error) VALUES(?,?,?) ON CONFLICT(task_id) DO UPDATE SET last_run_ns=excluded.last_run_ns,last_error=excluded.last_error`, id, started.UnixNano(), msg)
return e
}
func (s *localState) SeenURL(ctx context.Context, task, urlv string) bool {
var n int
_ = s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM seen WHERE task_id=? AND url=?`, task, urlv).Scan(&n)
return n > 0
}
func (s *localState) SeenHash(ctx context.Context, task, urlv, sha string) bool {
var n int
_ = s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM seen WHERE task_id=? AND url=? AND content_sha256=?`, task, urlv, sha).Scan(&n)
return n > 0
}
func (s *localState) MarkSeen(ctx context.Context, task, urlv, sha string) error {
_, err := s.db.ExecContext(ctx, `INSERT OR IGNORE INTO seen(task_id,url,content_sha256,seen_at_ns) VALUES(?,?,?,?)`, task, urlv, sha, time.Now().UTC().UnixNano())
return err
}
+91
View File
@@ -0,0 +1,91 @@
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)
}
}
+698
View File
@@ -0,0 +1,698 @@
package sourceagent
import (
"context"
"crypto/rand"
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"hash/fnv"
"math/bits"
"net/url"
"path/filepath"
"sort"
"strings"
"sync"
"time"
"unicode"
_ "modernc.org/sqlite"
)
type Store struct {
db *sql.DB
mu sync.Mutex
}
func OpenStore(dataDir string) (*Store, error) {
path := filepath.Join(dataDir, "source-agents.db")
db, err := sql.Open("sqlite", "file:"+filepath.ToSlash(path)+"?_pragma=busy_timeout(5000)&_pragma=foreign_keys(1)&_pragma=journal_mode(WAL)")
if err != nil {
return nil, err
}
db.SetMaxOpenConns(4)
store := &Store{db: db}
if err := store.init(context.Background()); err != nil {
_ = db.Close()
return nil, err
}
return store, nil
}
func (s *Store) Close() error {
if s == nil || s.db == nil {
return nil
}
return s.db.Close()
}
func (s *Store) init(ctx context.Context) error {
statements := []string{
`CREATE TABLE IF NOT EXISTS source_agents (
id TEXT PRIMARY KEY, name TEXT NOT NULL, token_hash TEXT NOT NULL, enabled INTEGER NOT NULL DEFAULT 1,
created_at_ns INTEGER NOT NULL, updated_at_ns INTEGER NOT NULL, last_seen_ns INTEGER NOT NULL DEFAULT 0,
last_error TEXT NOT NULL DEFAULT '', version TEXT NOT NULL DEFAULT ''
) WITHOUT ROWID`,
`CREATE TABLE IF NOT EXISTS source_tasks (
id TEXT PRIMARY KEY, agent_id TEXT NOT NULL REFERENCES source_agents(id) ON DELETE CASCADE,
name TEXT NOT NULL, type TEXT NOT NULL, url TEXT NOT NULL, enabled INTEGER NOT NULL DEFAULT 1,
poll_interval TEXT NOT NULL DEFAULT '4h', categories_json TEXT NOT NULL DEFAULT '[]', max_items INTEGER NOT NULL DEFAULT 20,
config_json TEXT NOT NULL DEFAULT '{}', created_at_ns INTEGER NOT NULL, updated_at_ns INTEGER NOT NULL
) WITHOUT ROWID`,
`CREATE INDEX IF NOT EXISTS idx_source_tasks_agent ON source_tasks(agent_id, enabled, updated_at_ns)`,
`CREATE TABLE IF NOT EXISTS source_inbox (
id TEXT PRIMARY KEY, agent_id TEXT NOT NULL, task_id TEXT NOT NULL, external_id TEXT NOT NULL DEFAULT '',
url TEXT NOT NULL, canonical_url TEXT NOT NULL, title TEXT NOT NULL, published_at_ns INTEGER NOT NULL DEFAULT 0,
discovered_at_ns INTEGER NOT NULL DEFAULT 0, received_at_ns INTEGER NOT NULL, updated_at_ns INTEGER NOT NULL,
language TEXT NOT NULL DEFAULT '', content_type TEXT NOT NULL DEFAULT '', text_content TEXT NOT NULL,
content_sha256 TEXT NOT NULL, source_name TEXT NOT NULL DEFAULT '', source_base_url TEXT NOT NULL DEFAULT '',
categories_json TEXT NOT NULL DEFAULT '[]', metadata_json TEXT NOT NULL DEFAULT '{}', signature INTEGER NOT NULL DEFAULT 0,
status TEXT NOT NULL DEFAULT 'received', relevance REAL NOT NULL DEFAULT 0, matched_node_id TEXT NOT NULL DEFAULT '',
UNIQUE(agent_id, task_id, canonical_url, content_sha256)
) WITHOUT ROWID`,
`CREATE INDEX IF NOT EXISTS idx_source_inbox_status_received ON source_inbox(status, received_at_ns)`,
`CREATE INDEX IF NOT EXISTS idx_source_inbox_candidate_updated ON source_inbox(status, updated_at_ns DESC)`,
}
for _, statement := range statements {
if _, err := s.db.ExecContext(ctx, statement); err != nil {
return err
}
}
_, _ = s.db.ExecContext(ctx, `UPDATE source_inbox SET status='received' WHERE status='processing' AND updated_at_ns<?`, time.Now().UTC().Add(-10*time.Minute).UnixNano())
return nil
}
func GenerateToken() (string, string, error) {
buf := make([]byte, 32)
if _, err := rand.Read(buf); err != nil {
return "", "", err
}
token := "brain_agent_" + hex.EncodeToString(buf)
return token, hashToken(token), nil
}
func hashToken(token string) string {
sum := sha256.Sum256([]byte(strings.TrimSpace(token)))
return hex.EncodeToString(sum[:])
}
func randomID(prefix string) string {
buf := make([]byte, 12)
_, _ = rand.Read(buf)
return prefix + "-" + hex.EncodeToString(buf)
}
func normalizeID(value, prefix string) string {
value = strings.ToLower(strings.TrimSpace(value))
var b strings.Builder
for _, r := range value {
if unicode.IsLetter(r) || unicode.IsDigit(r) || r == '-' || r == '_' {
b.WriteRune(r)
}
}
value = strings.Trim(b.String(), "-_")
if value == "" {
return randomID(prefix)
}
return value
}
func (s *Store) CreateAgent(ctx context.Context, id, name string) (Agent, string, error) {
id = normalizeID(id, "agent")
name = strings.TrimSpace(name)
if name == "" {
name = id
}
token, tokenHash, err := GenerateToken()
if err != nil {
return Agent{}, "", err
}
now := time.Now().UTC()
_, err = s.db.ExecContext(ctx, `INSERT INTO source_agents(id,name,token_hash,enabled,created_at_ns,updated_at_ns) VALUES(?,?,?,?,?,?)`, id, name, tokenHash, 1, now.UnixNano(), now.UnixNano())
if err != nil {
return Agent{}, "", err
}
return Agent{ID: id, Name: name, Enabled: true, CreatedAt: now, UpdatedAt: now}, token, nil
}
func (s *Store) RotateToken(ctx context.Context, id string) (string, error) {
token, tokenHash, err := GenerateToken()
if err != nil {
return "", err
}
result, err := s.db.ExecContext(ctx, `UPDATE source_agents SET token_hash=?, updated_at_ns=? WHERE id=?`, tokenHash, time.Now().UTC().UnixNano(), id)
if err != nil {
return "", err
}
n, _ := result.RowsAffected()
if n == 0 {
return "", sql.ErrNoRows
}
return token, nil
}
func (s *Store) SetAgentEnabled(ctx context.Context, id string, enabled bool) error {
v := 0
if enabled {
v = 1
}
result, err := s.db.ExecContext(ctx, `UPDATE source_agents SET enabled=?,updated_at_ns=? WHERE id=?`, v, time.Now().UTC().UnixNano(), id)
if err != nil {
return err
}
n, _ := result.RowsAffected()
if n == 0 {
return sql.ErrNoRows
}
return nil
}
func (s *Store) DeleteAgent(ctx context.Context, id string) error {
_, err := s.db.ExecContext(ctx, `DELETE FROM source_agents WHERE id=?`, id)
return err
}
func (s *Store) ListAgents(ctx context.Context) ([]Agent, error) {
rows, err := s.db.QueryContext(ctx, `SELECT a.id,a.name,a.enabled,a.created_at_ns,a.updated_at_ns,a.last_seen_ns,a.last_error,a.version,(SELECT COUNT(*) FROM source_tasks t WHERE t.agent_id=a.id) FROM source_agents a ORDER BY a.name,a.id`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Agent
for rows.Next() {
var a Agent
var en int
var c, u, seen int64
if err := rows.Scan(&a.ID, &a.Name, &en, &c, &u, &seen, &a.LastError, &a.Version, &a.TaskCount); err != nil {
return nil, err
}
a.Enabled = en != 0
a.CreatedAt = nsTime(c)
a.UpdatedAt = nsTime(u)
a.LastSeen = nsTime(seen)
out = append(out, a)
}
return out, rows.Err()
}
func (s *Store) Authenticate(ctx context.Context, token string) (Agent, error) {
if strings.TrimSpace(token) == "" {
return Agent{}, errors.New("empty agent token")
}
var a Agent
var en int
var c, u, seen int64
err := s.db.QueryRowContext(ctx, `SELECT id,name,enabled,created_at_ns,updated_at_ns,last_seen_ns,last_error,version FROM source_agents WHERE token_hash=?`, hashToken(token)).Scan(&a.ID, &a.Name, &en, &c, &u, &seen, &a.LastError, &a.Version)
if err != nil {
return Agent{}, err
}
a.Enabled = en != 0
if !a.Enabled {
return Agent{}, errors.New("agent disabled")
}
a.CreatedAt = nsTime(c)
a.UpdatedAt = nsTime(u)
a.LastSeen = nsTime(seen)
return a, nil
}
func validateTask(t Task) (Task, error) {
t.ID = normalizeID(t.ID, "task")
t.Name = strings.TrimSpace(t.Name)
if t.Name == "" {
t.Name = t.ID
}
t.Type = strings.ToLower(strings.TrimSpace(t.Type))
switch t.Type {
case "rss", "atom", "sitemap", "web":
default:
return t, fmt.Errorf("unsupported source task type %q", t.Type)
}
u, err := url.Parse(strings.TrimSpace(t.URL))
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") {
return t, fmt.Errorf("task URL must be absolute http(s)")
}
if u.User != nil {
return t, fmt.Errorf("task URL must not contain userinfo")
}
t.URL = u.String()
if strings.TrimSpace(t.PollInterval) == "" {
t.PollInterval = "4h"
}
d, err := time.ParseDuration(t.PollInterval)
if err != nil || d < 5*time.Minute || d > 30*24*time.Hour {
return t, fmt.Errorf("poll_interval must be between 5m and 720h")
}
if t.MaxItems <= 0 {
t.MaxItems = 20
}
if t.MaxItems > 500 {
return t, fmt.Errorf("max_items must be <=500")
}
if t.Config == nil {
t.Config = map[string]string{}
}
return t, nil
}
func (s *Store) UpsertTask(ctx context.Context, t Task) (Task, error) {
var err error
t, err = validateTask(t)
if err != nil {
return Task{}, err
}
if strings.TrimSpace(t.AgentID) == "" {
return Task{}, errors.New("agent_id required")
}
cats, _ := json.Marshal(uniqueStrings(t.Categories))
cfg, _ := json.Marshal(t.Config)
now := time.Now().UTC()
_, err = s.db.ExecContext(ctx, `INSERT INTO source_tasks(id,agent_id,name,type,url,enabled,poll_interval,categories_json,max_items,config_json,created_at_ns,updated_at_ns) VALUES(?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(id) DO UPDATE SET agent_id=excluded.agent_id,name=excluded.name,type=excluded.type,url=excluded.url,enabled=excluded.enabled,poll_interval=excluded.poll_interval,categories_json=excluded.categories_json,max_items=excluded.max_items,config_json=excluded.config_json,updated_at_ns=excluded.updated_at_ns`, t.ID, t.AgentID, t.Name, t.Type, t.URL, boolInt(t.Enabled), t.PollInterval, string(cats), t.MaxItems, string(cfg), now.UnixNano(), now.UnixNano())
if err != nil {
return Task{}, err
}
t.UpdatedAt = now
if t.CreatedAt.IsZero() {
t.CreatedAt = now
}
return t, nil
}
func (s *Store) DeleteTask(ctx context.Context, id string) error {
_, err := s.db.ExecContext(ctx, `DELETE FROM source_tasks WHERE id=?`, id)
return err
}
func (s *Store) ListTasks(ctx context.Context, agentID string) ([]Task, error) {
rows, err := s.db.QueryContext(ctx, `SELECT id,agent_id,name,type,url,enabled,poll_interval,categories_json,max_items,config_json,created_at_ns,updated_at_ns FROM source_tasks WHERE (?='' OR agent_id=?) ORDER BY name,id`, agentID, agentID)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Task
for rows.Next() {
var t Task
var en int
var cats, cfg string
var c, u int64
if err := rows.Scan(&t.ID, &t.AgentID, &t.Name, &t.Type, &t.URL, &en, &t.PollInterval, &cats, &t.MaxItems, &cfg, &c, &u); err != nil {
return nil, err
}
t.Enabled = en != 0
_ = json.Unmarshal([]byte(cats), &t.Categories)
_ = json.Unmarshal([]byte(cfg), &t.Config)
t.CreatedAt = nsTime(c)
t.UpdatedAt = nsTime(u)
out = append(out, t)
}
return out, rows.Err()
}
func (s *Store) RemoteConfig(ctx context.Context, agent Agent) (RemoteConfig, error) {
tasks, err := s.ListTasks(ctx, agent.ID)
if err != nil {
return RemoteConfig{}, err
}
return RemoteConfig{SchemaVersion: SchemaVersion, Agent: agent, Tasks: tasks, IssuedAt: time.Now().UTC()}, nil
}
func (s *Store) Heartbeat(ctx context.Context, agentID string, h Heartbeat) error {
_, err := s.db.ExecContext(ctx, `UPDATE source_agents SET last_seen_ns=?,last_error=?,version=?,updated_at_ns=? WHERE id=?`, time.Now().UTC().UnixNano(), strings.TrimSpace(h.LastError), strings.TrimSpace(h.Version), time.Now().UTC().UnixNano(), agentID)
return err
}
func normalizeDocument(d Document) (Document, error) {
d.URL = strings.TrimSpace(d.URL)
if d.URL == "" {
return d, errors.New("document url required")
}
u, err := url.Parse(d.URL)
if err != nil || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") {
return d, errors.New("invalid document url")
}
d.CanonicalURL = strings.TrimSpace(d.CanonicalURL)
if d.CanonicalURL == "" {
d.CanonicalURL = d.URL
}
d.Title = strings.TrimSpace(d.Title)
d.Text = strings.TrimSpace(d.Text)
if len([]rune(d.Text)) < 80 {
return d, errors.New("document text too short")
}
if d.DiscoveredAt.IsZero() {
d.DiscoveredAt = time.Now().UTC()
}
if d.ExternalID == "" {
d.ExternalID = d.CanonicalURL
}
d.Categories = uniqueStrings(d.Categories)
if d.Metadata == nil {
d.Metadata = map[string]any{}
}
if d.ContentSHA256 == "" {
sum := sha256.Sum256([]byte(d.Text))
d.ContentSHA256 = hex.EncodeToString(sum[:])
}
return d, nil
}
func (s *Store) Ingest(ctx context.Context, agentID, taskID string, docs []Document) (IngestResult, error) {
var taskOwner string
if err := s.db.QueryRowContext(ctx, `SELECT agent_id FROM source_tasks WHERE id=? AND enabled=1`, taskID).Scan(&taskOwner); err != nil {
return IngestResult{}, fmt.Errorf("unknown or disabled source task: %w", err)
}
if taskOwner != agentID {
return IngestResult{}, errors.New("source task does not belong to authenticated agent")
}
result := IngestResult{BatchID: randomID("batch")}
now := time.Now().UTC()
s.mu.Lock()
defer s.mu.Unlock()
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return result, err
}
defer tx.Rollback()
for _, raw := range docs {
d, e := normalizeDocument(raw)
if e != nil {
result.Rejected++
continue
}
cats, _ := json.Marshal(d.Categories)
meta, _ := json.Marshal(d.Metadata)
idSum := sha256.Sum256([]byte(agentID + "\x00" + taskID + "\x00" + d.CanonicalURL + "\x00" + d.ContentSHA256))
id := hex.EncodeToString(idSum[:12])
sig := int64(simhash64(d.Title + "\n" + d.Text))
res, e := tx.ExecContext(ctx, `INSERT OR IGNORE INTO source_inbox(id,agent_id,task_id,external_id,url,canonical_url,title,published_at_ns,discovered_at_ns,received_at_ns,updated_at_ns,language,content_type,text_content,content_sha256,source_name,source_base_url,categories_json,metadata_json,signature,status) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, id, agentID, taskID, d.ExternalID, d.URL, d.CanonicalURL, d.Title, timeNS(d.PublishedAt), timeNS(d.DiscoveredAt), now.UnixNano(), now.UnixNano(), d.Language, d.ContentType, d.Text, d.ContentSHA256, d.SourceName, d.SourceBaseURL, string(cats), string(meta), sig, "received")
if e != nil {
result.Rejected++
continue
}
n, _ := res.RowsAffected()
if n == 0 {
result.Duplicates++
} else {
result.Accepted++
}
}
if err := tx.Commit(); err != nil {
return result, err
}
return result, nil
}
func (s *Store) ClaimInbox(ctx context.Context, limit int) ([]InboxDocument, error) {
if limit < 1 {
limit = 1
}
if limit > 100 {
limit = 100
}
s.mu.Lock()
defer s.mu.Unlock()
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return nil, err
}
defer tx.Rollback()
rows, err := tx.QueryContext(ctx, `SELECT id FROM source_inbox WHERE status='received' ORDER BY received_at_ns LIMIT ?`, limit)
if err != nil {
return nil, err
}
var ids []string
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
rows.Close()
return nil, err
}
ids = append(ids, id)
}
rows.Close()
if len(ids) == 0 {
return nil, tx.Commit()
}
now := time.Now().UTC().UnixNano()
for _, id := range ids {
if _, err := tx.ExecContext(ctx, `UPDATE source_inbox SET status='processing',updated_at_ns=? WHERE id=? AND status='received'`, now, id); err != nil {
return nil, err
}
}
if err := tx.Commit(); err != nil {
return nil, err
}
return s.GetInboxByIDs(ctx, ids)
}
func (s *Store) CompleteClassification(ctx context.Context, id, status string, relevance float64, matchedNodeID string, meta map[string]any) error {
if status != "candidate" && status != "archived" {
return errors.New("invalid inbox classification status")
}
data, _ := json.Marshal(meta)
_, err := s.db.ExecContext(ctx, `UPDATE source_inbox SET status=?,relevance=?,matched_node_id=?,metadata_json=?,updated_at_ns=? WHERE id=?`, status, relevance, matchedNodeID, string(data), time.Now().UTC().UnixNano(), id)
return err
}
func (s *Store) MarkUsed(ctx context.Context, canonicalURLs []string) error {
now := time.Now().UTC().UnixNano()
for _, u := range uniqueStrings(canonicalURLs) {
if _, err := s.db.ExecContext(ctx, `UPDATE source_inbox SET status='used',updated_at_ns=? WHERE canonical_url=?`, now, u); err != nil {
return err
}
}
return nil
}
func (s *Store) GetInboxByIDs(ctx context.Context, ids []string) ([]InboxDocument, error) {
if len(ids) == 0 {
return nil, nil
}
out := make([]InboxDocument, 0, len(ids))
for _, id := range ids {
doc, err := s.getInbox(ctx, id)
if err == nil {
out = append(out, doc)
}
}
return out, nil
}
func (s *Store) getInbox(ctx context.Context, id string) (InboxDocument, error) {
row := s.db.QueryRowContext(ctx, `SELECT id,agent_id,task_id,external_id,url,canonical_url,title,published_at_ns,discovered_at_ns,received_at_ns,updated_at_ns,language,content_type,text_content,content_sha256,source_name,source_base_url,categories_json,metadata_json,status,relevance,matched_node_id FROM source_inbox WHERE id=?`, id)
return scanInbox(row)
}
type rowScanner interface{ Scan(...any) error }
func scanInbox(row rowScanner) (InboxDocument, error) {
var d InboxDocument
var ext, urlv, canon, title, lang, ctype, text, sha, source, base, cats, meta, status, matched string
var pub, disc, recv, upd int64
var rel float64
err := row.Scan(&d.ID, &d.AgentID, &d.TaskID, &ext, &urlv, &canon, &title, &pub, &disc, &recv, &upd, &lang, &ctype, &text, &sha, &source, &base, &cats, &meta, &status, &rel, &matched)
if err != nil {
return d, err
}
d.Document = Document{ExternalID: ext, URL: urlv, CanonicalURL: canon, Title: title, PublishedAt: nsTime(pub), DiscoveredAt: nsTime(disc), Language: lang, ContentType: ctype, Text: text, ContentSHA256: sha, SourceName: source, SourceBaseURL: base}
_ = json.Unmarshal([]byte(cats), &d.Document.Categories)
_ = json.Unmarshal([]byte(meta), &d.Metadata)
d.Status = status
d.Relevance = rel
d.MatchedNodeID = matched
d.ReceivedAt = nsTime(recv)
d.UpdatedAt = nsTime(upd)
return d, nil
}
func (s *Store) ListInbox(ctx context.Context, status string, limit int) ([]InboxDocument, error) {
if limit < 1 {
limit = 100
}
if limit > 1000 {
limit = 1000
}
rows, err := s.db.QueryContext(ctx, `SELECT id,agent_id,task_id,external_id,url,canonical_url,title,published_at_ns,discovered_at_ns,received_at_ns,updated_at_ns,language,content_type,text_content,content_sha256,source_name,source_base_url,categories_json,metadata_json,status,relevance,matched_node_id FROM source_inbox WHERE (?='' OR status=?) ORDER BY received_at_ns DESC LIMIT ?`, status, status, limit)
if err != nil {
return nil, err
}
defer rows.Close()
var out []InboxDocument
for rows.Next() {
d, err := scanInbox(rows)
if err != nil {
return nil, err
}
out = append(out, d)
}
return out, rows.Err()
}
func (s *Store) Stats(ctx context.Context) (InboxStats, error) {
rows, err := s.db.QueryContext(ctx, `SELECT status,COUNT(*) FROM source_inbox GROUP BY status`)
if err != nil {
return InboxStats{}, err
}
defer rows.Close()
var st InboxStats
for rows.Next() {
var status string
var n int
if err := rows.Scan(&status, &n); err != nil {
return st, err
}
st.Total += n
switch status {
case "received":
st.Received = n
case "processing":
st.Processing = n
case "candidate":
st.Candidate = n
case "archived":
st.Archived = n
case "used":
st.Used = n
}
}
return st, rows.Err()
}
func (s *Store) SearchCandidates(ctx context.Context, query string, limit int, maxAge time.Duration) ([]ScoredDocument, error) {
if limit < 1 {
limit = 3
}
if limit > 50 {
limit = 50
}
since := int64(0)
if maxAge > 0 {
since = time.Now().UTC().Add(-maxAge).UnixNano()
}
rows, err := s.db.QueryContext(ctx, `SELECT id,agent_id,task_id,external_id,url,canonical_url,title,published_at_ns,discovered_at_ns,received_at_ns,updated_at_ns,language,content_type,text_content,content_sha256,source_name,source_base_url,categories_json,metadata_json,status,relevance,matched_node_id,signature FROM source_inbox WHERE status IN ('candidate','used') AND (?=0 OR COALESCE(NULLIF(published_at_ns,0),received_at_ns)>=?) ORDER BY updated_at_ns DESC LIMIT 5000`, since, since)
if err != nil {
return nil, err
}
defer rows.Close()
qsig := simhash64(query)
qterms := termSet(query)
var out []ScoredDocument
for rows.Next() {
var d InboxDocument
var ext, urlv, canon, title, lang, ctype, text, sha, source, base, cats, meta, status, matched string
var pub, disc, recv, upd, sig int64
var rel float64
if err := rows.Scan(&d.ID, &d.AgentID, &d.TaskID, &ext, &urlv, &canon, &title, &pub, &disc, &recv, &upd, &lang, &ctype, &text, &sha, &source, &base, &cats, &meta, &status, &rel, &matched, &sig); err != nil {
return nil, err
}
d.Document = Document{ExternalID: ext, URL: urlv, CanonicalURL: canon, Title: title, PublishedAt: nsTime(pub), DiscoveredAt: nsTime(disc), Language: lang, ContentType: ctype, Text: text, ContentSHA256: sha, SourceName: source, SourceBaseURL: base}
_ = json.Unmarshal([]byte(cats), &d.Document.Categories)
_ = json.Unmarshal([]byte(meta), &d.Metadata)
d.Status = status
d.Relevance = rel
d.MatchedNodeID = matched
d.ReceivedAt = nsTime(recv)
d.UpdatedAt = nsTime(upd)
hamming := bits.OnesCount64(qsig ^ uint64(sig))
hashScore := 1 - float64(hamming)/64.0
lex := jaccard(qterms, termSet(title+" "+text))
score := 0.45*hashScore + 0.45*lex + 0.10*rel
if lex < 0.03 && hashScore < 0.58 {
continue
}
out = append(out, ScoredDocument{InboxDocument: d, QueryScore: score})
}
sort.Slice(out, func(i, j int) bool { return out[i].QueryScore > out[j].QueryScore })
if len(out) > limit {
out = out[:limit]
}
return out, nil
}
func boolInt(v bool) int {
if v {
return 1
}
return 0
}
func timeNS(t time.Time) int64 {
if t.IsZero() {
return 0
}
return t.UTC().UnixNano()
}
func nsTime(v int64) time.Time {
if v <= 0 {
return time.Time{}
}
return time.Unix(0, v).UTC()
}
func uniqueStrings(values []string) []string {
seen := map[string]bool{}
out := make([]string, 0, len(values))
for _, v := range values {
v = strings.TrimSpace(v)
if v == "" || seen[v] {
continue
}
seen[v] = true
out = append(out, v)
}
return out
}
func termSet(value string) map[string]bool {
out := map[string]bool{}
for _, r := range strings.FieldsFunc(strings.ToLower(value), func(r rune) bool { return !unicode.IsLetter(r) && !unicode.IsDigit(r) }) {
if len([]rune(r)) >= 3 {
out[r] = true
}
}
return out
}
func jaccard(a, b map[string]bool) float64 {
if len(a) == 0 || len(b) == 0 {
return 0
}
inter := 0
union := len(a)
for k := range b {
if a[k] {
inter++
} else {
union++
}
}
if union == 0 {
return 0
}
return float64(inter) / float64(union)
}
func simhash64(value string) uint64 {
weights := [64]int{}
for token := range termSet(value) {
h := fnv.New64a()
_, _ = h.Write([]byte(token))
x := h.Sum64()
for i := 0; i < 64; i++ {
if x&(1<<i) != 0 {
weights[i]++
} else {
weights[i]--
}
}
}
var sig uint64
for i, w := range weights {
if w >= 0 {
sig |= 1 << i
}
}
return sig
}
func (s *Store) ReleaseInbox(ctx context.Context, id, message string) error {
meta, _ := json.Marshal(map[string]any{"classification_error": strings.TrimSpace(message)})
_, err := s.db.ExecContext(ctx, `UPDATE source_inbox SET status='received',metadata_json=?,updated_at_ns=? WHERE id=?`, string(meta), time.Now().UTC().UnixNano(), id)
return err
}
+108
View File
@@ -0,0 +1,108 @@
package sourceagent
import "time"
const SchemaVersion = 1
type Agent struct {
ID string `json:"id"`
Name string `json:"name"`
Enabled bool `json:"enabled"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
LastSeen time.Time `json:"last_seen,omitempty"`
LastError string `json:"last_error,omitempty"`
Version string `json:"version,omitempty"`
TaskCount int `json:"task_count,omitempty"`
}
type Task struct {
ID string `json:"id"`
AgentID string `json:"agent_id"`
Name string `json:"name"`
Type string `json:"type"`
URL string `json:"url"`
Enabled bool `json:"enabled"`
PollInterval string `json:"poll_interval"`
Categories []string `json:"categories,omitempty"`
MaxItems int `json:"max_items,omitempty"`
Config map[string]string `json:"config,omitempty"`
CreatedAt time.Time `json:"created_at,omitempty"`
UpdatedAt time.Time `json:"updated_at,omitempty"`
}
type RemoteConfig struct {
SchemaVersion int `json:"schema_version"`
Agent Agent `json:"agent"`
Tasks []Task `json:"tasks"`
IssuedAt time.Time `json:"issued_at"`
}
type Document struct {
ExternalID string `json:"external_id,omitempty"`
URL string `json:"url"`
CanonicalURL string `json:"canonical_url,omitempty"`
Title string `json:"title"`
PublishedAt time.Time `json:"published_at,omitempty"`
DiscoveredAt time.Time `json:"discovered_at,omitempty"`
Language string `json:"language,omitempty"`
ContentType string `json:"content_type,omitempty"`
Text string `json:"text"`
ContentSHA256 string `json:"content_sha256,omitempty"`
SourceName string `json:"source_name,omitempty"`
SourceBaseURL string `json:"source_base_url,omitempty"`
Categories []string `json:"categories,omitempty"`
Metadata map[string]any `json:"metadata,omitempty"`
}
type IngestBatch struct {
SchemaVersion int `json:"schema_version"`
AgentID string `json:"agent_id,omitempty"`
TaskID string `json:"task_id"`
Documents []Document `json:"documents"`
}
type IngestResult struct {
Accepted int `json:"accepted"`
Duplicates int `json:"duplicates"`
Rejected int `json:"rejected"`
BatchID string `json:"batch_id"`
}
type Heartbeat struct {
AgentID string `json:"agent_id,omitempty"`
Version string `json:"version,omitempty"`
Status string `json:"status,omitempty"`
LastRunAt time.Time `json:"last_run_at,omitempty"`
LastError string `json:"last_error,omitempty"`
TasksChecked int `json:"tasks_checked,omitempty"`
Documents int `json:"documents,omitempty"`
Metadata map[string]any `json:"metadata,omitempty"`
}
type InboxDocument struct {
ID string `json:"id"`
AgentID string `json:"agent_id"`
TaskID string `json:"task_id"`
Document Document `json:"document"`
Status string `json:"status"`
Relevance float64 `json:"relevance,omitempty"`
MatchedNodeID string `json:"matched_node_id,omitempty"`
ReceivedAt time.Time `json:"received_at"`
UpdatedAt time.Time `json:"updated_at"`
Metadata map[string]any `json:"metadata,omitempty"`
}
type InboxStats struct {
Total int `json:"total"`
Received int `json:"received"`
Processing int `json:"processing"`
Candidate int `json:"candidate"`
Archived int `json:"archived"`
Used int `json:"used"`
}
type ScoredDocument struct {
InboxDocument
QueryScore float64 `json:"query_score"`
}