Files
jbergner 94dbd4ccab
All checks were successful
release-tag / release-image (push) Successful in 2m32s
RC-4
2026-08-09 18:41:47 +02:00

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