+131
-16
@@ -21,6 +21,7 @@ import (
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/local/glpi-neural-brain/internal/research"
|
||||
@@ -63,6 +64,7 @@ type Runner struct {
|
||||
controller *DockerController
|
||||
controllerMu sync.RWMutex
|
||||
controllerStatus DockerControllerStatus
|
||||
computeBusy atomic.Bool
|
||||
}
|
||||
|
||||
type runnerDiagnostics struct {
|
||||
@@ -255,6 +257,7 @@ func (r *Runner) Status() map[string]any {
|
||||
"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,
|
||||
@@ -267,6 +270,10 @@ func (r *Runner) Status() map[string]any {
|
||||
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)
|
||||
@@ -288,7 +295,7 @@ func (r *Runner) loop(ctx context.Context) {
|
||||
}
|
||||
r.runDue(ctx)
|
||||
if r.cfg.ComputeEnabled {
|
||||
r.runComputeOnce(ctx)
|
||||
r.triggerCompute(ctx)
|
||||
}
|
||||
if r.cfg.DockerControllerEnabled {
|
||||
r.runControllerOnce(ctx)
|
||||
@@ -301,6 +308,17 @@ func (r *Runner) loop(ctx context.Context) {
|
||||
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)
|
||||
@@ -309,7 +327,7 @@ func (r *Runner) loop(ctx context.Context) {
|
||||
r.runDue(ctx)
|
||||
case <-compute.C:
|
||||
if r.cfg.ComputeEnabled {
|
||||
r.runComputeOnce(ctx)
|
||||
r.triggerCompute(ctx)
|
||||
}
|
||||
case <-controller.C:
|
||||
if r.cfg.DockerControllerEnabled {
|
||||
@@ -364,12 +382,50 @@ func (r *Runner) refreshConfig(ctx context.Context) error {
|
||||
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()
|
||||
@@ -394,7 +450,7 @@ func (r *Runner) runDue(ctx context.Context) {
|
||||
if len(tasks) == 0 {
|
||||
return
|
||||
}
|
||||
sem := make(chan struct{}, r.cfg.Concurrency)
|
||||
sem := make(chan struct{}, r.effectiveConcurrency())
|
||||
var wg sync.WaitGroup
|
||||
for _, task := range tasks {
|
||||
task := task
|
||||
@@ -456,11 +512,64 @@ func (r *Runner) runTask(ctx context.Context, task Task) {
|
||||
_ = r.sendHeartbeat(ctx, h)
|
||||
}
|
||||
|
||||
func (r *Runner) runComputeOnce(ctx context.Context) {
|
||||
if r.runArticleQualityComputeOnce(ctx) {
|
||||
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
|
||||
}
|
||||
r.runVectorGraphComputeOnce(ctx)
|
||||
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 {
|
||||
@@ -468,13 +577,13 @@ func (r *Runner) runArticleQualityComputeOnce(ctx context.Context) bool {
|
||||
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 true
|
||||
return false
|
||||
}
|
||||
r.auth(req)
|
||||
resp, err := r.computeHTTP.Do(req)
|
||||
if err != nil {
|
||||
r.recordComputeResult(started, 0, err)
|
||||
return true
|
||||
return false
|
||||
}
|
||||
if resp.StatusCode == http.StatusNoContent {
|
||||
resp.Body.Close()
|
||||
@@ -485,14 +594,14 @@ func (r *Runner) runArticleQualityComputeOnce(ctx context.Context) bool {
|
||||
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 true
|
||||
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 true
|
||||
return false
|
||||
}
|
||||
result := ExecuteArticleQualityJob(job)
|
||||
result.AgentID = r.cfg.AgentID
|
||||
@@ -529,35 +638,35 @@ func (r *Runner) runArticleQualityComputeOnce(ctx context.Context) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (r *Runner) runVectorGraphComputeOnce(ctx context.Context) {
|
||||
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
|
||||
return false
|
||||
}
|
||||
r.auth(req)
|
||||
resp, err := r.computeHTTP.Do(req)
|
||||
if err != nil {
|
||||
r.recordComputeResult(started, 0, err)
|
||||
return
|
||||
return false
|
||||
}
|
||||
if resp.StatusCode == http.StatusNoContent {
|
||||
resp.Body.Close()
|
||||
return
|
||||
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
|
||||
return false
|
||||
}
|
||||
job, err := ReadVectorGraphJob(resp.Body, r.cfg.ComputeMaxBytes)
|
||||
resp.Body.Close()
|
||||
if err != nil {
|
||||
r.recordComputeResult(started, 0, err)
|
||||
return
|
||||
return false
|
||||
}
|
||||
result := ExecuteVectorGraphJob(job)
|
||||
result.AgentID = r.cfg.AgentID
|
||||
@@ -591,6 +700,7 @@ func (r *Runner) runVectorGraphComputeOnce(ctx context.Context) {
|
||||
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) {
|
||||
@@ -882,6 +992,11 @@ 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}
|
||||
|
||||
@@ -222,3 +222,17 @@ func TestProactiveSecurityInboxLifecycle(t *testing.T) {
|
||||
t.Fatalf("materialized document must remain searchable before SearXNG: %#v %v", results, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEffectiveConcurrencyUsesRemoteSpeedLimit(t *testing.T) {
|
||||
r := &Runner{
|
||||
cfg: RunnerConfig{Concurrency: 3},
|
||||
remote: RemoteConfig{Performance: AgentPerformance{SpeedMode: true, CPUWorkers: 12}},
|
||||
}
|
||||
if got := r.effectiveConcurrency(); got != 12 {
|
||||
t.Fatalf("speed concurrency=%d want 12", got)
|
||||
}
|
||||
r.remote.Performance.SpeedMode = false
|
||||
if got := r.effectiveConcurrency(); got != 3 {
|
||||
t.Fatalf("normal concurrency=%d want 3", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -33,11 +33,17 @@ type Task struct {
|
||||
UpdatedAt time.Time `json:"updated_at,omitempty"`
|
||||
}
|
||||
|
||||
type AgentPerformance struct {
|
||||
SpeedMode bool `json:"speed_mode"`
|
||||
CPUWorkers int `json:"cpu_tasks"`
|
||||
}
|
||||
|
||||
type RemoteConfig struct {
|
||||
SchemaVersion int `json:"schema_version"`
|
||||
Agent Agent `json:"agent"`
|
||||
Tasks []Task `json:"tasks"`
|
||||
IssuedAt time.Time `json:"issued_at"`
|
||||
SchemaVersion int `json:"schema_version"`
|
||||
Agent Agent `json:"agent"`
|
||||
Tasks []Task `json:"tasks"`
|
||||
Performance AgentPerformance `json:"performance"`
|
||||
IssuedAt time.Time `json:"issued_at"`
|
||||
}
|
||||
|
||||
type Document struct {
|
||||
|
||||
Reference in New Issue
Block a user