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" "sync/atomic" "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 ComputeEnabled bool ComputePollInterval time.Duration ComputeMaxBytes int64 DockerControllerEnabled bool DockerSocket string DockerComposeBinary string ControllerPollInterval time.Duration ControllerMaxDuration time.Duration } type Runner struct { cfg RunnerConfig http *http.Client computeHTTP *http.Client sourceHTTP *http.Client state *localState mu sync.RWMutex remote RemoteConfig wake chan struct{} diagMu sync.RWMutex diag runnerDiagnostics controller *DockerController controllerMu sync.RWMutex controllerStatus DockerControllerStatus computeBusy atomic.Bool } type runnerDiagnostics struct { StartedAt time.Time ConfigSource string CacheLoadedAt time.Time LastConfigAttemptAt time.Time LastConfigSuccessAt time.Time LastConfigError string LastHeartbeatAt time.Time LastHeartbeatSuccess time.Time LastHeartbeatError string LastConnectionErrorAt time.Time LastConnectionError string LastTaskRunAt time.Time LastTaskError string LastComputeRunAt time.Time LastComputeDurationMS int64 LastComputeError string ComputeCompleted uint64 LastControllerRunAt time.Time LastControllerDurationMS int64 LastControllerError string ControllerCompleted uint64 } 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 } if cfg.ComputePollInterval < time.Second { cfg.ComputePollInterval = 5 * time.Second } if cfg.ComputeMaxBytes < 8<<20 { cfg.ComputeMaxBytes = 128 << 20 } if cfg.ControllerPollInterval < time.Second { cfg.ControllerPollInterval = 5 * time.Second } if cfg.ControllerMaxDuration < 5*time.Second { cfg.ControllerMaxDuration = 15 * time.Minute } if strings.TrimSpace(cfg.DockerSocket) == "" { cfg.DockerSocket = "/var/run/docker.sock" } if strings.TrimSpace(cfg.DockerComposeBinary) == "" { cfg.DockerComposeBinary = "docker" } state, err := openLocalState(cfg.DataDir) if err != nil { return nil, err } r := &Runner{ cfg: cfg, http: &http.Client{Timeout: cfg.HTTPTimeout}, computeHTTP: &http.Client{Timeout: 10 * time.Minute}, sourceHTTP: research.NewSafeHTTPClient(cfg.AllowPrivate, cfg.HTTPTimeout), state: state, wake: make(chan struct{}, 1), } r.diag.StartedAt = time.Now().UTC() if cfg.DockerControllerEnabled { probeCtx, cancel := context.WithTimeout(context.Background(), 8*time.Second) controller, controllerErr := NewDockerController(probeCtx, cfg.DockerSocket, cfg.DockerComposeBinary) cancel() if controllerErr != nil { r.controllerStatus = DockerControllerStatus{Enabled: true, Socket: cfg.DockerSocket, Reachable: false, LastRefresh: time.Now().UTC(), LastError: controllerErr.Error()} } else { r.controller = controller statusCtx, statusCancel := context.WithTimeout(context.Background(), 8*time.Second) r.controllerStatus = controller.Status(statusCtx) statusCancel() } } return r, 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() r.diagMu.RLock() d := r.diag r.diagMu.RUnlock() loopbackURL := false if u, err := url.Parse(r.cfg.BrainURL); err == nil { host := strings.ToLower(u.Hostname()) loopbackURL = host == "localhost" || host == "127.0.0.1" || host == "::1" } lastSuccess := d.LastConfigSuccessAt if d.LastHeartbeatSuccess.After(lastSuccess) { lastSuccess = d.LastHeartbeatSuccess } connected := !lastSuccess.IsZero() && (d.LastConnectionErrorAt.IsZero() || !d.LastConnectionErrorAt.After(lastSuccess)) state := "never_connected" if connected { state = "connected" if d.LastConfigError != "" || d.LastHeartbeatError != "" { state = "degraded" } } else if !lastSuccess.IsZero() { state = "disconnected" } else if d.ConfigSource == "cache" { state = "cached_config_only" } warning := "" // Loopback is perfectly valid for a native same-host test. Only surface the // Docker-specific warning when the Agent is not actually connected. if loopbackURL && !connected { warning = "Brain URL points to loopback. In separate Docker containers, localhost/127.0.0.1 means the Agent container, not the Brain container. Use a shared service name or host.docker.internal:." } hint := "" errText := strings.ToLower(d.LastConnectionError) switch { case strings.Contains(errText, "http 401") || strings.Contains(errText, "unauthorized"): hint = "The Brain rejected this token. The Agent may no longer be registered in that Brain, may be disabled, or the token was rotated. Create/verify the Agent in the Brain UI and update BRAIN_AGENT_TOKEN." case strings.Contains(errText, "http 404"): hint = "The configured URL does not expose the Brain Agent API. Verify BRAIN_AGENT_BRAIN_URL and the published Brain port." case warning != "": hint = "For separate Docker containers, do not use 127.0.0.1/localhost as the Brain URL." } return map[string]any{ "ok": true, "mode": "agent", "agent_id": r.cfg.AgentID, "brain_url": r.cfg.BrainURL, "brain_connected": connected, "connection_state": state, "brain_url_warning": warning, "connection_hint": hint, "configured_tasks": len(remote.Tasks), "config_source": d.ConfigSource, "cache_loaded_at": d.CacheLoadedAt, "config_issued_at": remote.IssuedAt, "version": r.cfg.Version, "started_at": d.StartedAt, "last_config_attempt_at": d.LastConfigAttemptAt, "last_config_success_at": d.LastConfigSuccessAt, "last_config_error": d.LastConfigError, "last_heartbeat_at": d.LastHeartbeatAt, "last_heartbeat_success_at": d.LastHeartbeatSuccess, "last_heartbeat_error": d.LastHeartbeatError, "last_connection_error_at": d.LastConnectionErrorAt, "last_connection_error": d.LastConnectionError, "last_task_run_at": d.LastTaskRunAt, "last_task_error": d.LastTaskError, "compute_enabled": r.cfg.ComputeEnabled, "compute_poll_interval": r.cfg.ComputePollInterval, "speed_mode": remote.Performance.SpeedMode, "speed_cpu_tasks": r.effectiveConcurrency(), "last_compute_run_at": d.LastComputeRunAt, "last_compute_duration_ms": d.LastComputeDurationMS, "last_compute_error": d.LastComputeError, "compute_completed": d.ComputeCompleted, "docker_controller_enabled": r.cfg.DockerControllerEnabled, "last_controller_run_at": d.LastControllerRunAt, "last_controller_duration_ms": d.LastControllerDurationMS, "last_controller_error": d.LastControllerError, "controller_completed": d.ControllerCompleted, "docker_controller": r.currentControllerStatus(), } } func (r *Runner) loop(ctx context.Context) { refresh := time.NewTicker(r.cfg.ConfigRefresh) defer refresh.Stop() performance := time.NewTicker(5 * time.Second) defer performance.Stop() computeFast := time.NewTicker(250 * time.Millisecond) defer computeFast.Stop() run := time.NewTicker(30 * time.Second) defer run.Stop() heartbeat := time.NewTicker(time.Minute) defer heartbeat.Stop() compute := time.NewTicker(r.cfg.ComputePollInterval) defer compute.Stop() controller := time.NewTicker(r.cfg.ControllerPollInterval) defer controller.Stop() // Heartbeat independently of config loading. This lets the Brain see the // Agent even when its task configuration is temporarily broken. if err := r.sendHeartbeat(ctx, Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "starting"}); err != nil { slog.Warn("source agent initial heartbeat failed", "brain_url", r.cfg.BrainURL, "error", err) } if err := r.refreshConfig(ctx); err != nil { slog.Warn("source agent config refresh failed", "brain_url", r.cfg.BrainURL, "error", err) } else { slog.Info("source agent connected to brain", "brain_url", r.cfg.BrainURL, "tasks", len(r.currentTasks())) } r.runDue(ctx) if r.cfg.ComputeEnabled { r.triggerCompute(ctx) } if r.cfg.DockerControllerEnabled { r.runControllerOnce(ctx) } for { select { case <-ctx.Done(): return case <-refresh.C: if err := r.refreshConfig(ctx); err != nil { slog.Warn("source agent config refresh failed", "brain_url", r.cfg.BrainURL, "error", err) } case <-performance.C: if err := r.refreshPerformance(ctx); err != nil { slog.Debug("source agent performance refresh failed", "error", err) } case <-computeFast.C: r.mu.RLock() speed := r.remote.Performance.SpeedMode r.mu.RUnlock() if speed && r.cfg.ComputeEnabled { r.triggerCompute(ctx) } case <-heartbeat.C: if err := r.sendHeartbeat(ctx, Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "online", Metadata: map[string]any{"configured_tasks": len(r.currentTasks())}}); err != nil { slog.Warn("source agent heartbeat failed", "brain_url", r.cfg.BrainURL, "error", err) } case <-run.C: r.runDue(ctx) case <-compute.C: if r.cfg.ComputeEnabled { r.triggerCompute(ctx) } case <-controller.C: if r.cfg.DockerControllerEnabled { r.runControllerOnce(ctx) } case <-r.wake: r.runDue(ctx) } } } func (r *Runner) currentTasks() []Task { r.mu.RLock() defer r.mu.RUnlock() return append([]Task(nil), r.remote.Tasks...) } func (r *Runner) refreshConfig(ctx context.Context) error { r.diagMu.Lock() r.diag.LastConfigAttemptAt = time.Now().UTC() r.diagMu.Unlock() req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/config", nil) if err != nil { r.recordConfigResult(err) return err } r.auth(req) resp, err := r.http.Do(req) if err != nil { r.recordConfigResult(err) return err } defer resp.Body.Close() if resp.StatusCode/100 != 2 { b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) err := fmt.Errorf("brain config HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b))) r.recordConfigResult(err) return err } var cfg RemoteConfig if err := json.NewDecoder(io.LimitReader(resp.Body, 4<<20)).Decode(&cfg); err != nil { r.recordConfigResult(err) return err } if cfg.Agent.ID != "" && cfg.Agent.ID != r.cfg.AgentID { err := fmt.Errorf("brain returned config for unexpected agent %q", cfg.Agent.ID) r.recordConfigResult(err) return err } r.mu.Lock() r.remote = cfg r.mu.Unlock() r.saveCachedConfig(cfg) r.recordConfigResult(nil) if cfg.Performance.SpeedMode && r.cfg.ComputeEnabled { r.triggerCompute(ctx) } if err := r.sendHeartbeat(ctx, Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "online", Metadata: map[string]any{"configured_tasks": len(cfg.Tasks)}}); err != nil { slog.Warn("source agent heartbeat after config refresh failed", "error", err) } return nil } func (r *Runner) refreshPerformance(ctx context.Context) error { req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/performance", 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 { b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) return fmt.Errorf("brain performance HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b))) } var perf AgentPerformance if err := json.NewDecoder(io.LimitReader(resp.Body, 64<<10)).Decode(&perf); err != nil { return err } if perf.CPUWorkers < 1 { perf.CPUWorkers = 1 } if perf.CPUWorkers > 256 { perf.CPUWorkers = 256 } r.mu.Lock() previous := r.remote.Performance r.remote.Performance = perf r.mu.Unlock() if perf.SpeedMode && (!previous.SpeedMode || perf.CPUWorkers != previous.CPUWorkers) && r.cfg.ComputeEnabled { r.triggerCompute(ctx) } return nil } func (r *Runner) recordConfigResult(err error) { r.diagMu.Lock() defer r.diagMu.Unlock() now := time.Now().UTC() if err != nil { r.diag.LastConfigError = err.Error() r.diag.LastConnectionErrorAt = now r.diag.LastConnectionError = err.Error() return } r.diag.LastConfigSuccessAt = now r.diag.LastConfigError = "" r.diag.ConfigSource = "brain" r.diag.LastConnectionErrorAt = time.Time{} r.diag.LastConnectionError = "" } 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.effectiveConcurrency()) 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() r.diagMu.Lock() r.diag.LastTaskRunAt = started r.diagMu.Unlock() 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) r.diagMu.Lock() if err != nil { r.diag.LastTaskError = err.Error() } else { r.diag.LastTaskError = "" } r.diagMu.Unlock() 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) effectiveConcurrency() int { r.mu.RLock() perf := r.remote.Performance r.mu.RUnlock() limit := r.cfg.Concurrency if perf.SpeedMode && perf.CPUWorkers > limit { limit = perf.CPUWorkers } if limit < 1 { limit = 1 } if limit > 256 { limit = 256 } return limit } func (r *Runner) triggerCompute(ctx context.Context) { if !r.computeBusy.CompareAndSwap(false, true) { return } go func() { defer r.computeBusy.Store(false) r.runComputeDispatch(ctx) }() } func (r *Runner) runComputeDispatch(ctx context.Context) { r.mu.RLock() speed := r.remote.Performance.SpeedMode r.mu.RUnlock() workers := 1 if speed { workers = r.effectiveConcurrency() } var wg sync.WaitGroup for i := 0; i < workers; i++ { wg.Add(1) go func() { defer wg.Done() for { if ctx.Err() != nil || !r.runComputeOnce(ctx) { return } if !speed { return } } }() } wg.Wait() } func (r *Runner) runComputeOnce(ctx context.Context) bool { if r.runArticleQualityComputeOnce(ctx) { return true } return r.runVectorGraphComputeOnce(ctx) } func (r *Runner) runArticleQualityComputeOnce(ctx context.Context) bool { started := time.Now().UTC() req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/compute/article-quality/claim", nil) if err != nil { r.recordComputeResult(started, 0, err) return false } r.auth(req) resp, err := r.computeHTTP.Do(req) if err != nil { r.recordComputeResult(started, 0, err) return false } if resp.StatusCode == http.StatusNoContent { resp.Body.Close() return false } if resp.StatusCode/100 != 2 { body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) resp.Body.Close() err = fmt.Errorf("brain article-quality claim HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) r.recordComputeResult(started, 0, err) return false } var job ArticleQualityComputeRequest err = json.NewDecoder(io.LimitReader(resp.Body, 8<<20)).Decode(&job) resp.Body.Close() if err != nil { r.recordComputeResult(started, 0, err) return false } result := ExecuteArticleQualityJob(job) result.AgentID = r.cfg.AgentID data, err := json.Marshal(result) if err == nil { resultCtx, cancel := context.WithTimeout(ctx, 2*time.Minute) var resultReq *http.Request resultReq, err = http.NewRequestWithContext(resultCtx, http.MethodPost, r.cfg.BrainURL+"/api/v1/agent/compute/article-quality/"+url.PathEscape(job.JobID)+"/result", bytes.NewReader(data)) if err == nil { resultReq.Header.Set("Content-Type", "application/json") r.auth(resultReq) var resultResp *http.Response resultResp, err = r.computeHTTP.Do(resultReq) if err == nil { defer resultResp.Body.Close() if resultResp.StatusCode/100 != 2 { body, _ := io.ReadAll(io.LimitReader(resultResp.Body, 4096)) err = fmt.Errorf("brain article-quality result HTTP %d: %s", resultResp.StatusCode, strings.TrimSpace(string(body))) } } } cancel() } r.recordComputeResult(started, result.DurationMS, err) h := Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "compute", LastRunAt: started, Metadata: map[string]any{"compute_job_id": job.JobID, "compute_kind": ComputeKindArticleQuality, "compute_duration_ms": result.DurationMS, "article_quality_score": result.Quality.Score, "article_quality_passed": result.Quality.Passed, "no_model_call": true}} if err != nil { h.Status = "error" h.LastError = err.Error() slog.Warn("source agent article quality job failed", "job_id", job.JobID, "error", err) } else { slog.Info("source agent article quality job completed", "job_id", job.JobID, "duration_ms", result.DurationMS, "score", result.Quality.Score) } _ = r.sendHeartbeat(ctx, h) return true } func (r *Runner) runVectorGraphComputeOnce(ctx context.Context) bool { started := time.Now().UTC() req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/compute/claim", nil) if err != nil { r.recordComputeResult(started, 0, err) return false } r.auth(req) resp, err := r.computeHTTP.Do(req) if err != nil { r.recordComputeResult(started, 0, err) return false } if resp.StatusCode == http.StatusNoContent { resp.Body.Close() return false } if resp.StatusCode/100 != 2 { body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) resp.Body.Close() err = fmt.Errorf("brain compute claim HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) r.recordComputeResult(started, 0, err) return false } job, err := ReadVectorGraphJob(resp.Body, r.cfg.ComputeMaxBytes) resp.Body.Close() if err != nil { r.recordComputeResult(started, 0, err) return false } result := ExecuteVectorGraphJob(job) result.AgentID = r.cfg.AgentID data, err := json.Marshal(result) if err == nil { resultCtx, cancel := context.WithTimeout(ctx, 2*time.Minute) var resultReq *http.Request resultReq, err = http.NewRequestWithContext(resultCtx, http.MethodPost, r.cfg.BrainURL+"/api/v1/agent/compute/"+url.PathEscape(job.Header.JobID)+"/result", bytes.NewReader(data)) if err == nil { resultReq.Header.Set("Content-Type", "application/json") r.auth(resultReq) var resultResp *http.Response resultResp, err = r.computeHTTP.Do(resultReq) if err == nil { defer resultResp.Body.Close() if resultResp.StatusCode/100 != 2 { body, _ := io.ReadAll(io.LimitReader(resultResp.Body, 4096)) err = fmt.Errorf("brain compute result HTTP %d: %s", resultResp.StatusCode, strings.TrimSpace(string(body))) } } } cancel() } r.recordComputeResult(started, result.DurationMS, err) h := Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "compute", LastRunAt: started, Metadata: map[string]any{"compute_job_id": job.Header.JobID, "compute_kind": ComputeKindVectorGraph, "compute_duration_ms": result.DurationMS, "compute_indexed": result.Primary.Stats.Indexed, "compute_links": result.Primary.Stats.Links + result.Orphan.Stats.Links}} if err != nil { h.Status = "error" h.LastError = err.Error() slog.Warn("source agent compute job failed", "job_id", job.Header.JobID, "error", err) } else { slog.Info("source agent compute job completed", "job_id", job.Header.JobID, "duration_ms", result.DurationMS, "links", result.Primary.Stats.Links+result.Orphan.Stats.Links) } _ = r.sendHeartbeat(ctx, h) return true } func (r *Runner) recordComputeResult(started time.Time, durationMS int64, err error) { r.diagMu.Lock() defer r.diagMu.Unlock() r.diag.LastComputeRunAt = started r.diag.LastComputeDurationMS = durationMS if err != nil { r.diag.LastComputeError = err.Error() return } r.diag.LastComputeError = "" if durationMS > 0 { r.diag.ComputeCompleted++ } } 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)]*href\s*=\s*["']([^"'#]+)["'][^>]*>(.*?)`) 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, "security_proactive": strings.ToLower(strings.TrimSpace(task.Config["security_proactive"]))}}) 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 { if h.Metadata == nil { h.Metadata = map[string]any{} } r.mu.RLock() perf := r.remote.Performance r.mu.RUnlock() h.Metadata["speed_mode"] = perf.SpeedMode h.Metadata["speed_cpu_tasks"] = r.effectiveConcurrency() capabilities := []string{} if r.cfg.ComputeEnabled { h.Metadata["compute_kinds"] = []string{ComputeKindVectorGraph, ComputeKindArticleQuality} capabilities = append(capabilities, ComputeKindVectorGraph, ComputeKindArticleQuality) } else { h.Metadata["compute_kinds"] = []string{} } if r.cfg.DockerControllerEnabled { statusCtx, cancel := context.WithTimeout(ctx, 5*time.Second) status := r.refreshControllerStatus(statusCtx) cancel() h.Metadata["docker_controller"] = status if status.Reachable { capabilities = append(capabilities, CapabilityDockerController) if status.ComposeAvailable { capabilities = append(capabilities, CapabilityDockerCompose) } } } h.Metadata["capabilities"] = uniqueStrings(capabilities) r.diagMu.Lock() r.diag.LastHeartbeatAt = time.Now().UTC() r.diagMu.Unlock() data, _ := json.Marshal(h) req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.cfg.BrainURL+"/api/v1/agent/heartbeat", bytes.NewReader(data)) if err != nil { r.recordHeartbeatResult(err) return err } req.Header.Set("Content-Type", "application/json") r.auth(req) resp, err := r.http.Do(req) if err != nil { r.recordHeartbeatResult(err) return err } defer resp.Body.Close() if resp.StatusCode/100 != 2 { b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) err := fmt.Errorf("heartbeat HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b))) r.recordHeartbeatResult(err) return err } r.recordHeartbeatResult(nil) return nil } func (r *Runner) recordHeartbeatResult(err error) { r.diagMu.Lock() defer r.diagMu.Unlock() now := time.Now().UTC() if err != nil { r.diag.LastHeartbeatError = err.Error() r.diag.LastConnectionErrorAt = now r.diag.LastConnectionError = err.Error() return } r.diag.LastHeartbeatSuccess = now r.diag.LastHeartbeatError = "" r.diag.LastConnectionErrorAt = time.Time{} r.diag.LastConnectionError = "" } 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() r.diagMu.Lock() if r.diag.ConfigSource == "" { r.diag.ConfigSource = "cache" } r.diag.CacheLoadedAt = time.Now().UTC() r.diagMu.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 } func (r *Runner) currentControllerStatus() DockerControllerStatus { r.controllerMu.RLock() defer r.controllerMu.RUnlock() return r.controllerStatus } func (r *Runner) refreshControllerStatus(ctx context.Context) DockerControllerStatus { r.controllerMu.Lock() defer r.controllerMu.Unlock() if !r.cfg.DockerControllerEnabled { r.controllerStatus = DockerControllerStatus{Enabled: false} return r.controllerStatus } if r.controller == nil { controller, err := NewDockerController(ctx, r.cfg.DockerSocket, r.cfg.DockerComposeBinary) if err != nil { r.controllerStatus = DockerControllerStatus{Enabled: true, Socket: r.cfg.DockerSocket, Reachable: false, LastRefresh: time.Now().UTC(), LastError: err.Error()} return r.controllerStatus } r.controller = controller } r.controllerStatus = r.controller.Status(ctx) return r.controllerStatus } func (r *Runner) runControllerOnce(ctx context.Context) { if !r.cfg.DockerControllerEnabled { return } statusCtx, statusCancel := context.WithTimeout(ctx, 8*time.Second) status := r.refreshControllerStatus(statusCtx) statusCancel() if !status.Reachable || r.controller == nil { return } started := time.Now().UTC() req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/controller/claim", nil) if err != nil { r.recordControllerResult(started, 0, err) return } r.auth(req) resp, err := r.computeHTTP.Do(req) if err != nil { r.recordControllerResult(started, 0, err) return } if resp.StatusCode == http.StatusNoContent { resp.Body.Close() return } if resp.StatusCode/100 != 2 { body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) resp.Body.Close() r.recordControllerResult(started, 0, fmt.Errorf("brain controller claim HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body)))) return } var claim ControllerClaim err = json.NewDecoder(io.LimitReader(resp.Body, 2<<20)).Decode(&claim) resp.Body.Close() if err != nil { r.recordControllerResult(started, 0, err) return } if claim.SchemaVersion != SchemaVersion || claim.Job.ID == "" { r.recordControllerResult(started, 0, errors.New("invalid controller claim")) return } // MaxJobDuration is deliberately not serialized as a raw duration value. // Rehydrate the central Brain policy from its textual representation so the // host-side Agent cannot silently widen the centrally configured job limit. limit := r.cfg.ControllerMaxDuration if configured, parseErr := time.ParseDuration(strings.TrimSpace(claim.Policy.MaxJobDurationText)); parseErr == nil && configured > 0 && configured < limit { limit = configured } jobCtx, cancel := context.WithTimeout(ctx, limit) jobCtx, authCancel := r.controllerAuthorizedContext(jobCtx, claim.Job.ID) result := r.controller.Execute(jobCtx, claim) authCancel() cancel() result.AgentID = r.cfg.AgentID data, marshalErr := json.Marshal(result) if marshalErr == nil { resultCtx, resultCancel := context.WithTimeout(ctx, 30*time.Second) resultReq, reqErr := http.NewRequestWithContext(resultCtx, http.MethodPost, r.cfg.BrainURL+"/api/v1/agent/controller/"+url.PathEscape(claim.Job.ID)+"/result", bytes.NewReader(data)) if reqErr == nil { resultReq.Header.Set("Content-Type", "application/json") r.auth(resultReq) resultResp, doErr := r.computeHTTP.Do(resultReq) if doErr != nil { err = doErr } else { if resultResp.StatusCode/100 != 2 { body, _ := io.ReadAll(io.LimitReader(resultResp.Body, 4096)) err = fmt.Errorf("brain controller result HTTP %d: %s", resultResp.StatusCode, strings.TrimSpace(string(body))) } resultResp.Body.Close() } } else { err = reqErr } resultCancel() } else { err = marshalErr } if result.Error != "" && err == nil { err = errors.New(result.Error) } r.recordControllerResult(started, result.DurationMS, err) _ = r.sendHeartbeat(ctx, Heartbeat{AgentID: r.cfg.AgentID, Version: r.cfg.Version, Status: "controller", LastRunAt: started, LastError: result.Error, Metadata: map[string]any{"controller_job_id": claim.Job.ID, "controller_kind": claim.Job.Kind, "controller_status": result.Status}}) } func (r *Runner) controllerAuthorizedContext(parent context.Context, jobID string) (context.Context, context.CancelFunc) { ctx, cancel := context.WithCancel(parent) check := func() bool { req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.cfg.BrainURL+"/api/v1/agent/controller/"+url.PathEscape(jobID)+"/authorized", nil) if err != nil { return false } r.auth(req) resp, err := r.http.Do(req) if err != nil { return false } defer resp.Body.Close() return resp.StatusCode/100 == 2 } if !check() { cancel() return ctx, cancel } go func() { ticker := time.NewTicker(2 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if !check() { cancel() return } } } }() return ctx, cancel } func (r *Runner) recordControllerResult(started time.Time, durationMS int64, err error) { r.diagMu.Lock() defer r.diagMu.Unlock() r.diag.LastControllerRunAt = started if durationMS <= 0 { durationMS = time.Since(started).Milliseconds() } r.diag.LastControllerDurationMS = durationMS if err != nil { r.diag.LastControllerError = err.Error() return } r.diag.LastControllerError = "" r.diag.ControllerCompleted++ }