All checks were successful
release-tag / release-image (push) Successful in 2m32s
1356 lines
43 KiB
Go
1356 lines
43 KiB
Go
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:<published-port>."
|
|
}
|
|
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)<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, "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++
|
|
}
|