v1.6.0
release-tag / release-image (push) Successful in 6m50s

This commit is contained in:
2026-09-02 10:26:50 +02:00
parent 8c4ce2d6c2
commit 9a4370e4df
272 changed files with 31244 additions and 5980 deletions
+103
View File
@@ -254,6 +254,108 @@ func run() (retErr error) {
}
}
// Master/subagent orchestrator and graph bootstrap. Every option is explicit
// and environment-owned only when set, preserving admin-managed config otherwise.
workerEnv := []string{
"NEUROFORGE_WORKER_LEASE_SECONDS", "NEUROFORGE_WORKER_HEARTBEAT_SECONDS", "NEUROFORGE_WORKER_STALE_AFTER_SECONDS",
"NEUROFORGE_WORKER_DEFAULT_MAX_ATTEMPTS", "NEUROFORGE_WORKER_RETRY_BACKOFF_SECONDS", "NEUROFORGE_WORKER_MAX_QUEUED_JOBS",
"NEUROFORGE_WORKER_JOB_RETENTION_HOURS", "NEUROFORGE_WORKER_MAX_TERMINAL_JOBS",
"NEUROFORGE_WORKER_MASTER_APPLY_MAX_ATTEMPTS", "NEUROFORGE_WORKER_MASTER_APPLY_BACKOFF_SECONDS",
"NEUROFORGE_GRAPH_BACKFILL_ENABLED", "NEUROFORGE_GRAPH_BACKFILL_INTERVAL_SECONDS", "NEUROFORGE_GRAPH_BACKFILL_BATCH_SIZE",
"NEUROFORGE_GRAPH_BACKFILL_MAX_QUEUED", "NEUROFORGE_GRAPH_BACKFILL_MIN_DEGREE", "NEUROFORGE_GRAPH_CANDIDATE_MULTIPLIER",
"NEUROFORGE_GRAPH_RETRY_AFTER_MINUTES", "NEUROFORGE_GRAPH_REQUIRE_WORKER", "NEUROFORGE_GRAPH_MAX_HOPS",
"NEUROFORGE_GRAPH_HOP_DECAY", "NEUROFORGE_GRAPH_MAX_EXPANSION", "NEUROFORGE_GRAPH_MIN_EDGE_WEIGHT",
"NEUROFORGE_OFFLOAD_CHAT", "NEUROFORGE_OFFLOAD_EMBEDDINGS", "NEUROFORGE_DISTRIBUTED_INFERENCE_WAIT_SECONDS",
}
hasWorkerEnv := false
for _, name := range workerEnv {
if _, ok := os.LookupEnv(name); ok {
hasWorkerEnv = true
break
}
}
if hasWorkerEnv {
cfg := s.Config()
if v, ok := envInt("NEUROFORGE_WORKER_LEASE_SECONDS"); ok {
cfg.Worker.LeaseSeconds = v
}
if v, ok := envInt("NEUROFORGE_WORKER_HEARTBEAT_SECONDS"); ok {
cfg.Worker.HeartbeatSeconds = v
}
if v, ok := envInt("NEUROFORGE_WORKER_STALE_AFTER_SECONDS"); ok {
cfg.Worker.StaleAfterSeconds = v
}
if v, ok := envInt("NEUROFORGE_WORKER_DEFAULT_MAX_ATTEMPTS"); ok {
cfg.Worker.DefaultMaxAttempts = v
}
if v, ok := envInt("NEUROFORGE_WORKER_RETRY_BACKOFF_SECONDS"); ok {
cfg.Worker.RetryBackoffSeconds = v
}
if v, ok := envInt("NEUROFORGE_WORKER_MAX_QUEUED_JOBS"); ok {
cfg.Worker.MaxQueuedJobs = v
}
if v, ok := envInt("NEUROFORGE_WORKER_JOB_RETENTION_HOURS"); ok {
cfg.Worker.JobRetentionHours = v
}
if v, ok := envInt("NEUROFORGE_WORKER_MAX_TERMINAL_JOBS"); ok {
cfg.Worker.MaxTerminalJobs = v
}
if v, ok := envInt("NEUROFORGE_WORKER_MASTER_APPLY_MAX_ATTEMPTS"); ok {
cfg.Worker.MasterApplyMaxAttempts = v
}
if v, ok := envInt("NEUROFORGE_WORKER_MASTER_APPLY_BACKOFF_SECONDS"); ok {
cfg.Worker.MasterApplyBackoffSeconds = v
}
if v, ok := envBool("NEUROFORGE_GRAPH_BACKFILL_ENABLED"); ok {
cfg.Worker.GraphBackfillEnabled = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_BACKFILL_INTERVAL_SECONDS"); ok {
cfg.Worker.GraphBackfillIntervalS = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_BACKFILL_BATCH_SIZE"); ok {
cfg.Worker.GraphBackfillBatchSize = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_BACKFILL_MAX_QUEUED"); ok {
cfg.Worker.GraphBackfillMaxQueued = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_BACKFILL_MIN_DEGREE"); ok {
cfg.Worker.GraphBackfillMinDegree = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_CANDIDATE_MULTIPLIER"); ok {
cfg.Worker.GraphCandidateMultiplier = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_RETRY_AFTER_MINUTES"); ok {
cfg.Worker.GraphRetryAfterMinutes = v
}
if v, ok := envBool("NEUROFORGE_GRAPH_REQUIRE_WORKER"); ok {
cfg.Worker.RequireWorkerForGraph = v
}
if v, ok := envBool("NEUROFORGE_OFFLOAD_CHAT"); ok {
cfg.Worker.OffloadChat = v
}
if v, ok := envBool("NEUROFORGE_OFFLOAD_EMBEDDINGS"); ok {
cfg.Worker.OffloadEmbeddings = v
}
if v, ok := envInt("NEUROFORGE_DISTRIBUTED_INFERENCE_WAIT_SECONDS"); ok {
cfg.Worker.DistributedInferenceWaitS = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_MAX_HOPS"); ok {
cfg.Brain.GraphMaxHops = v
}
if v, ok := envFloat("NEUROFORGE_GRAPH_HOP_DECAY"); ok {
cfg.Brain.GraphHopDecay = v
}
if v, ok := envInt("NEUROFORGE_GRAPH_MAX_EXPANSION"); ok {
cfg.Brain.GraphMaxExpansion = v
}
if v, ok := envFloat("NEUROFORGE_GRAPH_MIN_EDGE_WEIGHT"); ok {
cfg.Brain.GraphMinEdgeWeight = v
}
if err := s.UpdateConfig(cfg); err != nil {
return fmt.Errorf("apply orchestrator/graph environment bootstrap: %w", err)
}
}
r := provider.NewRouter(s)
c := cost.New(s)
b := brain.New(s, r, c)
@@ -336,6 +438,7 @@ func run() (retErr error) {
maintenanceCtx, stopMaintenance := context.WithCancel(rootCtx)
defer stopMaintenance()
go b.RunV6Maintenance(maintenanceCtx)
go b.RunOrchestrator(maintenanceCtx)
api := httpapi.New(s, b, r, c)
if v, ok := envBool("NEUROFORGE_READINESS_OLLAMA_LIVE"); ok {
api.SetReadinessOllamaLive(v)
+421 -42
View File
@@ -2,6 +2,7 @@ package main
import (
"bytes"
"context"
"encoding/json"
"flag"
"fmt"
@@ -9,20 +10,30 @@ import (
"log"
"net/http"
"os"
"os/signal"
"sort"
"strconv"
"strings"
"sync"
"syscall"
"time"
"neuroforge/internal/core"
"neuroforge/internal/vector"
)
const workerVersion = "1.6.0"
type relinkPayload struct {
TargetID string `json:"target_id"`
Target []float32 `json:"target"`
Candidates []struct {
ID string `json:"id"`
Vector []float32 `json:"vector"`
TargetID string `json:"target_id"`
TargetVersion int64 `json:"target_version"`
TargetFingerprint string `json:"target_fingerprint"`
Target []float32 `json:"target"`
Candidates []struct {
ID string `json:"id"`
Version int64 `json:"version"`
Fingerprint string `json:"fingerprint"`
Vector []float32 `json:"vector"`
} `json:"candidates"`
K int `json:"k"`
MinSimilarity float64 `json:"min_similarity"`
@@ -36,58 +47,226 @@ type relinkResult struct {
Neighbors []neighbor `json:"neighbors"`
}
type modelEmbedPayload struct {
Text string `json:"text"`
Model string `json:"model,omitempty"`
}
type modelChatPayload struct {
Instructions string `json:"instructions,omitempty"`
Input string `json:"input"`
Model string `json:"model,omitempty"`
MaxOutput int `json:"max_output,omitempty"`
JSONMode bool `json:"json_mode,omitempty"`
}
type modelUsage struct {
InputTokens int64 `json:"input_tokens"`
CachedTokens int64 `json:"cached_tokens,omitempty"`
OutputTokens int64 `json:"output_tokens"`
}
type modelResult struct {
Text string `json:"text,omitempty"`
Vector []float32 `json:"vector,omitempty"`
Usage modelUsage `json:"usage"`
Provider string `json:"provider"`
Model string `json:"model"`
NodeID string `json:"node_id"`
}
type workerConfig struct {
ID string
Server string
Token string
Interval time.Duration
Heartbeat time.Duration
ResourceClass string
Capabilities []string
MaxConcurrency int
Hostname string
OllamaURL string
OllamaChatModel string
OllamaEmbedModel string
OllamaNumCtx int
OllamaKeepAlive string
}
type leaseSet struct {
mu sync.RWMutex
m map[string]string
}
func (l *leaseSet) add(id, token string) {
l.mu.Lock()
l.m[id] = token
l.mu.Unlock()
}
func (l *leaseSet) del(id string) {
l.mu.Lock()
delete(l.m, id)
l.mu.Unlock()
}
func (l *leaseSet) snapshot() map[string]string {
l.mu.RLock()
defer l.mu.RUnlock()
out := make(map[string]string, len(l.m))
for k, v := range l.m {
out[k] = v
}
return out
}
func main() {
server := flag.String("server", "http://localhost:8080", "NeuroForge server")
id := flag.String("id", hostname(), "worker id")
server := flag.String("server", "http://localhost:8080", "NeuroForge master/orchestrator server")
id := flag.String("id", hostname(), "subagent worker id")
interval := flag.Duration("interval", 2*time.Second, "poll interval")
flag.Parse()
token := strings.TrimSpace(os.Getenv("NEUROFORGE_WORKER_TOKEN"))
if token == "" {
log.Fatal("NEUROFORGE_WORKER_TOKEN is required")
}
client := &http.Client{Timeout: 180 * time.Second}
log.Printf("worker %s polling %s", *id, *server)
for {
job, err := claim(client, *server, token, *id)
resource := lowerDefault(os.Getenv("NEUROFORGE_WORKER_RESOURCE_CLASS"), "cpu")
caps := csv(os.Getenv("NEUROFORGE_WORKER_CAPABILITIES"))
if len(caps) == 0 {
if resource == "gpu" {
caps = []string{"gpu", "model.chat", "model.embed"}
} else {
caps = []string{"cpu", "vector.relink"}
}
}
maxConcurrency := envInt("NEUROFORGE_WORKER_MAX_CONCURRENCY", 1)
if maxConcurrency < 1 {
maxConcurrency = 1
}
heartbeat := envDuration("NEUROFORGE_WORKER_HEARTBEAT_INTERVAL", 15*time.Second)
cfg := workerConfig{
ID: *id, Server: strings.TrimRight(*server, "/"), Token: token, Interval: *interval,
Heartbeat: heartbeat, ResourceClass: resource, Capabilities: caps,
MaxConcurrency: maxConcurrency, Hostname: hostname(),
OllamaURL: strings.TrimRight(firstNonEmpty(os.Getenv("NEUROFORGE_WORKER_OLLAMA_URL"), os.Getenv("OLLAMA_BASE_URL"), os.Getenv("OLLAMA_URL")), "/"),
OllamaChatModel: firstNonEmpty(os.Getenv("NEUROFORGE_WORKER_OLLAMA_CHAT_MODEL"), os.Getenv("OLLAMA_MODEL")),
OllamaEmbedModel: firstNonEmpty(os.Getenv("NEUROFORGE_WORKER_OLLAMA_EMBEDDING_MODEL"), os.Getenv("OLLAMA_EMBEDDING_MODEL")),
OllamaNumCtx: envInt("NEUROFORGE_WORKER_OLLAMA_NUM_CTX", 8192),
OllamaKeepAlive: firstNonEmpty(os.Getenv("NEUROFORGE_WORKER_OLLAMA_KEEP_ALIVE"), "10m"),
}
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
// Control-plane requests must never hang forever on a dead master. Inference
// uses its own client and is bounded by the durable job timeout/context.
controlClient := &http.Client{Timeout: 30 * time.Second}
inferenceClient := &http.Client{Timeout: 0}
leases := &leaseSet{m: map[string]string{}}
if err := register(controlClient, cfg); err != nil {
log.Printf("initial register failed: %v", err)
}
go heartbeatLoop(ctx, controlClient, cfg, leases)
log.Printf("subagent %s polling %s resource=%s caps=%s concurrency=%d", cfg.ID, cfg.Server, cfg.ResourceClass, strings.Join(cfg.Capabilities, ","), cfg.MaxConcurrency)
var wg sync.WaitGroup
for i := 0; i < cfg.MaxConcurrency; i++ {
wg.Add(1)
go func(slot int) {
defer wg.Done()
pollLoop(ctx, controlClient, inferenceClient, cfg, leases, slot)
}(i)
}
<-ctx.Done()
log.Printf("subagent shutdown requested; waiting for active slots")
wg.Wait()
}
func waitOrDone(ctx context.Context, d time.Duration) bool {
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return false
case <-t.C:
return true
}
}
func pollLoop(root context.Context, controlClient, inferenceClient *http.Client, cfg workerConfig, leases *leaseSet, slot int) {
for root.Err() == nil {
job, err := claim(controlClient, cfg)
if err != nil {
log.Printf("claim: %v", err)
time.Sleep(*interval)
log.Printf("slot=%d claim: %v", slot, err)
if !waitOrDone(root, cfg.Interval) {
return
}
continue
}
if job == nil {
time.Sleep(*interval)
if !waitOrDone(root, cfg.Interval) {
return
}
continue
}
res, jobErr := run(job)
if err := complete(client, *server, token, *id, job.ID, res, jobErr); err != nil {
log.Printf("complete %s: %v", job.ID, err)
leases.add(job.ID, job.LeaseToken)
jobCtx := root
cancel := func() {}
if job.TimeoutSeconds > 0 {
jobCtx, cancel = context.WithTimeout(root, time.Duration(job.TimeoutSeconds)*time.Second)
}
res, jobErr := run(jobCtx, inferenceClient, cfg, job)
cancel()
if root.Err() != nil && jobErr == "" {
jobErr = root.Err().Error()
}
if err := complete(controlClient, cfg, job.ID, job.LeaseToken, res, jobErr); err != nil {
log.Printf("slot=%d complete %s: %v", slot, job.ID, err)
} else if jobErr != "" {
log.Printf("slot=%d job %s %s returned error: %s", slot, job.ID, job.Type, jobErr)
} else {
log.Printf("job %s %s done", job.ID, job.Type)
log.Printf("slot=%d job %s %s complete", slot, job.ID, job.Type)
}
leases.del(job.ID)
}
}
func workerBody(cfg workerConfig, leases map[string]string) map[string]any {
return map[string]any{
"worker_id": cfg.ID, "resource_class": cfg.ResourceClass, "capabilities": cfg.Capabilities,
"max_concurrency": cfg.MaxConcurrency, "version": workerVersion, "hostname": cfg.Hostname,
"active_leases": leases,
}
}
func register(c *http.Client, cfg workerConfig) error {
return postJSON(c, cfg.Server+"/api/v1/worker/register", cfg.Token, workerBody(cfg, nil), nil, http.StatusOK)
}
func heartbeatLoop(ctx context.Context, c *http.Client, cfg workerConfig, leases *leaseSet) {
t := time.NewTicker(cfg.Heartbeat)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
var state core.WorkerState
if err := postJSON(c, cfg.Server+"/api/v1/worker/heartbeat", cfg.Token, workerBody(cfg, leases.snapshot()), &state, http.StatusOK); err != nil {
log.Printf("heartbeat: %v", err)
}
}
}
}
func hostname() string {
h, _ := os.Hostname()
if h == "" {
h = "worker"
}
return h
}
func claim(c *http.Client, server, token, id string) (*core.Job, error) {
body, _ := json.Marshal(map[string]string{"worker_id": id})
req, _ := http.NewRequest("POST", strings.TrimRight(server, "/")+"/api/v1/worker/claim", bytes.NewReader(body))
req.Header.Set("Authorization", "Bearer "+token)
func claim(c *http.Client, cfg workerConfig) (*core.Job, error) {
body, _ := json.Marshal(workerBody(cfg, nil))
req, _ := http.NewRequest("POST", cfg.Server+"/api/v1/worker/claim", bytes.NewReader(body))
req.Header.Set("Authorization", "Bearer "+cfg.Token)
req.Header.Set("Content-Type", "application/json")
resp, err := c.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == 204 {
if resp.StatusCode == http.StatusNoContent {
return nil, nil
}
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 8<<20))
if resp.StatusCode != 200 {
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 32<<20))
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("HTTP %d: %s", resp.StatusCode, raw)
}
var j core.Job
@@ -96,7 +275,8 @@ func claim(c *http.Client, server, token, id string) (*core.Job, error) {
}
return &j, nil
}
func run(j *core.Job) (json.RawMessage, string) {
func run(ctx context.Context, c *http.Client, cfg workerConfig, j *core.Job) (json.RawMessage, string) {
switch j.Type {
case "vector.relink":
var p relinkPayload
@@ -104,10 +284,10 @@ func run(j *core.Job) (json.RawMessage, string) {
return nil, err.Error()
}
out := relinkResult{TargetID: p.TargetID}
for _, c := range p.Candidates {
sim := vector.Cosine(p.Target, c.Vector)
for _, candidate := range p.Candidates {
sim := vector.Cosine(p.Target, candidate.Vector)
if sim >= p.MinSimilarity {
out.Neighbors = append(out.Neighbors, neighbor{ID: c.ID, Similarity: sim})
out.Neighbors = append(out.Neighbors, neighbor{ID: candidate.ID, Similarity: sim})
}
}
sort.Slice(out.Neighbors, func(i, k int) bool { return out.Neighbors[i].Similarity > out.Neighbors[k].Similarity })
@@ -116,13 +296,134 @@ func run(j *core.Job) (json.RawMessage, string) {
}
b, _ := json.Marshal(out)
return b, ""
case "model.embed":
if !hasCap(cfg.Capabilities, "model.embed") {
return nil, "worker lacks model.embed capability"
}
var p modelEmbedPayload
if err := json.Unmarshal(j.Payload, &p); err != nil {
return nil, err.Error()
}
out, err := ollamaEmbed(ctx, c, cfg, p)
if err != nil {
return nil, err.Error()
}
b, _ := json.Marshal(out)
return b, ""
case "model.chat":
if !hasCap(cfg.Capabilities, "model.chat") {
return nil, "worker lacks model.chat capability"
}
var p modelChatPayload
if err := json.Unmarshal(j.Payload, &p); err != nil {
return nil, err.Error()
}
out, err := ollamaChat(ctx, c, cfg, p)
if err != nil {
return nil, err.Error()
}
b, _ := json.Marshal(out)
return b, ""
default:
return nil, "unsupported job type: " + j.Type
}
}
func complete(c *http.Client, server, token, id, jobID string, result json.RawMessage, jobErr string) error {
body, _ := json.Marshal(map[string]any{"worker_id": id, "job_id": jobID, "result": result, "error": jobErr})
req, _ := http.NewRequest("POST", strings.TrimRight(server, "/")+"/api/v1/worker/complete", bytes.NewReader(body))
func ollamaEmbed(ctx context.Context, c *http.Client, cfg workerConfig, p modelEmbedPayload) (modelResult, error) {
if cfg.OllamaURL == "" {
return modelResult{}, fmt.Errorf("NEUROFORGE_WORKER_OLLAMA_URL is required for model.embed")
}
model := firstNonEmpty(p.Model, cfg.OllamaEmbedModel)
if model == "" || strings.TrimSpace(p.Text) == "" {
return modelResult{}, fmt.Errorf("embedding model and text are required")
}
body := map[string]any{"model": model, "input": p.Text, "keep_alive": cfg.OllamaKeepAlive}
var resp struct {
Embeddings [][]float32 `json:"embeddings"`
PromptEvalCount int64 `json:"prompt_eval_count"`
}
if err := postOllama(ctx, c, cfg.OllamaURL+"/api/embed", body, &resp); err != nil {
return modelResult{}, err
}
if len(resp.Embeddings) == 0 || len(resp.Embeddings[0]) == 0 {
return modelResult{}, fmt.Errorf("ollama embedding response is empty")
}
return modelResult{Vector: resp.Embeddings[0], Usage: modelUsage{InputTokens: resp.PromptEvalCount}, Provider: "ollama", Model: model, NodeID: cfg.ID}, nil
}
func ollamaChat(ctx context.Context, c *http.Client, cfg workerConfig, p modelChatPayload) (modelResult, error) {
if cfg.OllamaURL == "" {
return modelResult{}, fmt.Errorf("NEUROFORGE_WORKER_OLLAMA_URL is required for model.chat")
}
model := firstNonEmpty(p.Model, cfg.OllamaChatModel)
if model == "" || strings.TrimSpace(p.Input) == "" {
return modelResult{}, fmt.Errorf("chat model and input are required")
}
messages := []map[string]string{}
if strings.TrimSpace(p.Instructions) != "" {
messages = append(messages, map[string]string{"role": "system", "content": p.Instructions})
}
messages = append(messages, map[string]string{"role": "user", "content": p.Input})
body := map[string]any{"model": model, "messages": messages, "stream": false, "keep_alive": cfg.OllamaKeepAlive}
if p.JSONMode {
body["format"] = "json"
}
options := map[string]any{}
if cfg.OllamaNumCtx > 0 {
options["num_ctx"] = cfg.OllamaNumCtx
}
if p.MaxOutput > 0 {
options["num_predict"] = p.MaxOutput
}
if len(options) > 0 {
body["options"] = options
}
var resp struct {
Message struct {
Content string `json:"content"`
} `json:"message"`
PromptEvalCount int64 `json:"prompt_eval_count"`
EvalCount int64 `json:"eval_count"`
}
if err := postOllama(ctx, c, cfg.OllamaURL+"/api/chat", body, &resp); err != nil {
return modelResult{}, err
}
return modelResult{Text: strings.TrimSpace(resp.Message.Content), Usage: modelUsage{InputTokens: resp.PromptEvalCount, OutputTokens: resp.EvalCount}, Provider: "ollama", Model: model, NodeID: cfg.ID}, nil
}
func postOllama(ctx context.Context, c *http.Client, url string, body any, out any) error {
raw, _ := json.Marshal(body)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(raw))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
b, _ := io.ReadAll(io.LimitReader(resp.Body, 32<<20))
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("ollama HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
}
if err := json.Unmarshal(b, out); err != nil {
return fmt.Errorf("decode ollama response: %w", err)
}
return nil
}
func complete(c *http.Client, cfg workerConfig, jobID, leaseToken string, result json.RawMessage, jobErr string) error {
body := map[string]any{"worker_id": cfg.ID, "job_id": jobID, "lease_token": leaseToken, "result": result, "error": jobErr}
return postJSON(c, cfg.Server+"/api/v1/worker/complete", cfg.Token, body, nil, http.StatusOK)
}
func postJSON(c *http.Client, url, token string, body any, out any, want int) error {
raw, _ := json.Marshal(body)
req, err := http.NewRequest(http.MethodPost, url, bytes.NewReader(raw))
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Content-Type", "application/json")
resp, err := c.Do(req)
@@ -130,9 +431,87 @@ func complete(c *http.Client, server, token, id, jobID string, result json.RawMe
return err
}
defer resp.Body.Close()
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 2<<20))
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("HTTP %d: %s", resp.StatusCode, raw)
b, _ := io.ReadAll(io.LimitReader(resp.Body, 8<<20))
if resp.StatusCode != want {
return fmt.Errorf("HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(b)))
}
if out != nil && len(bytes.TrimSpace(b)) > 0 {
if err := json.Unmarshal(b, out); err != nil {
return err
}
}
return nil
}
func hostname() string {
h, _ := os.Hostname()
if h == "" {
h = "worker"
}
return h
}
func csv(v string) []string {
seen := map[string]bool{}
out := []string{}
for _, x := range strings.Split(v, ",") {
x = strings.ToLower(strings.TrimSpace(x))
if x != "" && !seen[x] {
seen[x] = true
out = append(out, x)
}
}
sort.Strings(out)
return out
}
func hasCap(caps []string, want string) bool {
want = strings.ToLower(strings.TrimSpace(want))
for _, x := range caps {
if strings.ToLower(strings.TrimSpace(x)) == want {
return true
}
}
return false
}
func lowerDefault(v, def string) string {
v = strings.ToLower(strings.TrimSpace(v))
if v == "" {
return def
}
return v
}
func firstNonEmpty(xs ...string) string {
for _, x := range xs {
if strings.TrimSpace(x) != "" {
return strings.TrimSpace(x)
}
}
return ""
}
func envInt(name string, def int) int {
v := strings.TrimSpace(os.Getenv(name))
if v == "" {
return def
}
n, err := strconv.Atoi(v)
if err != nil {
return def
}
return n
}
func envDuration(name string, def time.Duration) time.Duration {
v := strings.TrimSpace(os.Getenv(name))
if v == "" {
return def
}
d, err := time.ParseDuration(v)
if err != nil || d <= 0 {
return def
}
return d
}
+30 -12
View File
@@ -29,28 +29,46 @@ services:
retries: 4
start_period: 10s
worker:
worker-cpu:
build:
context: .
target: worker
command:
- -server
- http://neuroforge:8080
- -id
- worker-compose-1
command: ["-server", "http://neuroforge:8080", "-id", "cpu-compose-1"]
environment:
NEUROFORGE_WORKER_TOKEN: ${NEUROFORGE_WORKER_TOKEN:?Set a unique NeuroForge worker token}
NEUROFORGE_WORKER_RESOURCE_CLASS: cpu
NEUROFORGE_WORKER_CAPABILITIES: cpu,vector.relink
NEUROFORGE_WORKER_MAX_CONCURRENCY: ${NEUROFORGE_CPU_WORKER_CONCURRENCY:-2}
depends_on:
neuroforge:
condition: service_healthy
restart: unless-stopped
read_only: true
tmpfs:
- /tmp:size=32m,mode=1777
security_opt:
- no-new-privileges:true
cap_drop:
- ALL
tmpfs: ["/tmp:size=32m,mode=1777"]
security_opt: ["no-new-privileges:true"]
cap_drop: ["ALL"]
worker-gpu:
build:
context: .
target: worker
command: ["-server", "http://neuroforge:8080", "-id", "gpu-compose-1"]
environment:
NEUROFORGE_WORKER_TOKEN: ${NEUROFORGE_WORKER_TOKEN:?Set a unique NeuroForge worker token}
NEUROFORGE_WORKER_RESOURCE_CLASS: gpu
NEUROFORGE_WORKER_CAPABILITIES: gpu,model.chat,model.embed
NEUROFORGE_WORKER_MAX_CONCURRENCY: ${NEUROFORGE_GPU_WORKER_CONCURRENCY:-1}
NEUROFORGE_WORKER_OLLAMA_URL: ${OLLAMA_BASE_URL:-http://ollama:11434}
NEUROFORGE_WORKER_OLLAMA_CHAT_MODEL: ${OLLAMA_MODEL:-gemma3}
NEUROFORGE_WORKER_OLLAMA_EMBEDDING_MODEL: ${OLLAMA_EMBEDDING_MODEL:-embeddinggemma}
depends_on:
neuroforge:
condition: service_healthy
restart: unless-stopped
read_only: true
tmpfs: ["/tmp:size=32m,mode=1777"]
security_opt: ["no-new-privileges:true"]
cap_drop: ["ALL"]
volumes:
neuroforge-data:
+103 -17
View File
@@ -3,10 +3,14 @@ package brain
import (
"bytes"
"context"
"crypto/sha256"
"encoding/binary"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"math"
"net/http"
"regexp"
"sort"
@@ -87,6 +91,12 @@ func (e *Engine) embed(ctx context.Context, text string) (provider.EmbedResult,
if route == "" {
route = "auto"
}
if (route == "auto" || route == "ollama") && nodeID == "" {
if res, used, offloadErr := e.distributedEmbed(ctx, model, text); used && offloadErr == nil {
costUSD, recErr := e.cost.Record(res.Provider, res.Model, "embedding", res.Usage)
return res, costUSD, recErr
}
}
if route == "auto" {
res, err := e.router.EmbedOn(ctx, "ollama", model, nodeID, text)
if err == nil {
@@ -160,6 +170,12 @@ func (e *Engine) chatModelLimitOnMode(ctx context.Context, providerName, model,
if route == "" {
route = "auto"
}
if (route == "auto" || route == "ollama") && nodeID == "" {
if res, used, offloadErr := e.distributedChat(ctx, model, instructions, input, maxOutput, jsonMode); used && offloadErr == nil {
costUSD, recErr := e.cost.Record(res.Provider, res.Model, "chat", res.Usage)
return res, costUSD, recErr
}
}
if route == "auto" {
var res provider.ChatResult
var err error
@@ -986,15 +1002,19 @@ func minFloat(a, b float64) float64 {
}
type relinkPayload struct {
TargetID string `json:"target_id"`
Target []float32 `json:"target"`
Candidates []relinkCandidate `json:"candidates"`
K int `json:"k"`
MinSimilarity float64 `json:"min_similarity"`
TargetID string `json:"target_id"`
TargetVersion int64 `json:"target_version"`
TargetFingerprint string `json:"target_fingerprint"`
Target []float32 `json:"target"`
Candidates []relinkCandidate `json:"candidates"`
K int `json:"k"`
MinSimilarity float64 `json:"min_similarity"`
}
type relinkCandidate struct {
ID string `json:"id"`
Vector []float32 `json:"vector"`
ID string `json:"id"`
Version int64 `json:"version"`
Fingerprint string `json:"fingerprint"`
Vector []float32 `json:"vector"`
}
type RelinkResult struct {
TargetID string `json:"target_id"`
@@ -1004,16 +1024,49 @@ type RelinkResult struct {
} `json:"neighbors"`
}
func vectorFingerprint(v []float32) string {
h := sha256.New()
var buf [4]byte
for _, x := range v {
binary.LittleEndian.PutUint32(buf[:], math.Float32bits(x))
_, _ = h.Write(buf[:])
}
sum := h.Sum(nil)
return hex.EncodeToString(sum[:12])
}
var errObsoleteRelink = errors.New("obsolete vector.relink result")
func (e *Engine) enqueueRelink(m *core.Memory) (*core.Job, error) {
cfg := e.store.Config()
snap := e.store.MemoriesSnapshot()
p := relinkPayload{TargetID: m.ID, Target: m.Vector, K: cfg.Brain.RecallK, MinSimilarity: cfg.Brain.MinSimilarity}
for _, x := range snap {
if x.ID != m.ID && len(x.Vector) == len(m.Vector) {
p.Candidates = append(p.Candidates, relinkCandidate{ID: x.ID, Vector: x.Vector})
}
if m == nil || m.ID == "" || len(m.Vector) == 0 {
return nil, errors.New("relink target requires id and vector")
}
return e.store.EnqueueJob("vector.relink", p)
mult := cfg.Worker.GraphCandidateMultiplier
if mult < 2 {
mult = 4
}
candidateK := cfg.Brain.RecallK * mult
if candidateK < cfg.Brain.RecallK+4 {
candidateK = cfg.Brain.RecallK + 4
}
if candidateK > 256 {
candidateK = 256
}
hits := e.store.SearchVector(m.Vector, candidateK+1, cfg.Brain.MinSimilarity, 0)
p := relinkPayload{TargetID: m.ID, TargetVersion: m.Version, TargetFingerprint: vectorFingerprint(m.Vector), Target: append([]float32(nil), m.Vector...), K: cfg.Brain.RecallK, MinSimilarity: cfg.Brain.MinSimilarity}
for _, h := range hits {
if h.Memory.ID == m.ID || len(h.Memory.Vector) != len(m.Vector) {
continue
}
p.Candidates = append(p.Candidates, relinkCandidate{ID: h.Memory.ID, Version: h.Memory.Version, Fingerprint: vectorFingerprint(h.Memory.Vector), Vector: append([]float32(nil), h.Memory.Vector...)})
}
return e.store.EnqueueJobSpec(store.JobSpec{
Type: "vector.relink", Payload: p, Priority: 20, ResourceClass: "cpu",
RequiredCapabilities: []string{"cpu", "vector.relink"},
RequiresMasterApply: true,
IdempotencyKey: "vector.relink:" + m.ID + ":v" + strconv.FormatInt(m.Version, 10),
})
}
func (e *Engine) localRelink(m *core.Memory) error {
cfg := e.store.Config()
@@ -1026,20 +1079,53 @@ func (e *Engine) localRelink(m *core.Memory) error {
return nil
}
func (e *Engine) ApplyJobResult(j *core.Job) error {
if j.Type != "vector.relink" || j.Status != "done" {
if j == nil || j.Type != "vector.relink" || (j.Status != "done" && j.Status != "apply_wait") {
return nil
}
var p relinkPayload
if err := json.Unmarshal(j.Payload, &p); err != nil {
return fmt.Errorf("decode relink payload during master apply: %w", err)
}
current, ok := e.store.GetMemory(p.TargetID)
if !ok || current.Version != p.TargetVersion || vectorFingerprint(current.Vector) != p.TargetFingerprint {
// The worker computed against an older delete/recreate or memory version.
// Treat the result as obsolete rather than poisoning the current graph; the
// backfill planner will schedule the current version again.
return errObsoleteRelink
}
var r RelinkResult
if err := json.Unmarshal(j.Result, &r); err != nil {
return err
}
if r.TargetID != p.TargetID {
return fmt.Errorf("relink result target %q does not match payload target %q", r.TargetID, p.TargetID)
}
candidateMeta := make(map[string]relinkCandidate, len(p.Candidates))
for _, c := range p.Candidates {
candidateMeta[c.ID] = c
}
sort.Slice(r.Neighbors, func(i, j int) bool { return r.Neighbors[i].Similarity > r.Neighbors[j].Similarity })
for _, n := range r.Neighbors {
if err := e.reinforcePair(r.TargetID, n.ID, n.Similarity, n.Similarity); err != nil {
meta, expected := candidateMeta[n.ID]
if !expected {
// A worker may only return neighbors from the bounded master-selected
// candidate set. Rejecting extras closes a trust-boundary gap.
continue
}
neighbor, ok := e.store.GetMemory(n.ID)
if !ok || neighbor.Version != meta.Version || vectorFingerprint(neighbor.Vector) != meta.Fingerprint {
continue
}
cfg := e.store.Config()
delta := cfg.Brain.LearningRate * n.Similarity
if delta <= 0 {
delta = n.Similarity
}
if err := e.store.UpsertRelationEvidence(r.TargetID, n.ID, "semantic_similarity", n.Similarity, delta, cfg.Brain.MaxSynapseWeight); err != nil {
return err
}
}
return nil
return e.store.MarkGraphLinked(r.TargetID)
}
// SearchByProvenanceSources embeds text once and searches only the requested
@@ -0,0 +1,107 @@
package brain
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"neuroforge/internal/provider"
"neuroforge/internal/store"
)
type distributedModelResult struct {
Text string `json:"text,omitempty"`
Vector []float32 `json:"vector,omitempty"`
Usage struct {
InputTokens int64 `json:"input_tokens"`
CachedTokens int64 `json:"cached_tokens,omitempty"`
OutputTokens int64 `json:"output_tokens"`
} `json:"usage"`
Provider string `json:"provider"`
Model string `json:"model"`
NodeID string `json:"node_id"`
}
func (e *Engine) waitDistributedJob(ctx context.Context, jID string) (*distributedModelResult, error) {
t := time.NewTicker(50 * time.Millisecond)
defer t.Stop()
for {
select {
case <-ctx.Done():
_ = e.store.CancelJob(jID, "caller context ended: "+ctx.Err().Error())
return nil, ctx.Err()
case <-t.C:
j, ok := e.store.Job(jID)
if !ok {
return nil, errors.New("distributed job disappeared")
}
switch j.Status {
case "done":
var out distributedModelResult
if err := json.Unmarshal(j.Result, &out); err != nil {
return nil, fmt.Errorf("decode distributed model result: %w", err)
}
return &out, nil
case "failed", "canceled":
return nil, fmt.Errorf("distributed job %s: %s", j.Status, j.Error)
}
}
}
}
func (e *Engine) distributedEmbed(ctx context.Context, model, text string) (provider.EmbedResult, bool, error) {
cfg := e.store.Config()
if !cfg.Worker.OffloadEmbeddings || !e.store.HasLiveWorker("gpu", "gpu", "model.embed") {
return provider.EmbedResult{}, false, nil
}
wait := cfg.Worker.DistributedInferenceWaitS
if wait <= 0 {
wait = 180
}
j, err := e.store.EnqueueJobSpec(store.JobSpec{
Type: "model.embed", Payload: map[string]any{"text": text, "model": model},
Priority: 100, ResourceClass: "gpu", RequiredCapabilities: []string{"gpu", "model.embed"},
MaxAttempts: 2, TimeoutSeconds: wait,
})
if err != nil {
return provider.EmbedResult{}, true, err
}
jobCtx, cancel := context.WithTimeout(ctx, time.Duration(wait)*time.Second)
defer cancel()
out, err := e.waitDistributedJob(jobCtx, j.ID)
if err != nil {
return provider.EmbedResult{}, true, err
}
if len(out.Vector) == 0 {
return provider.EmbedResult{}, true, errors.New("distributed embedding returned empty vector")
}
return provider.EmbedResult{Vector: out.Vector, Provider: out.Provider, Model: out.Model, NodeID: out.NodeID, Usage: provider.Usage{InputTokens: out.Usage.InputTokens, CachedTokens: out.Usage.CachedTokens, OutputTokens: out.Usage.OutputTokens}}, true, nil
}
func (e *Engine) distributedChat(ctx context.Context, model, instructions, input string, maxOutput int, jsonMode bool) (provider.ChatResult, bool, error) {
cfg := e.store.Config()
if !cfg.Worker.OffloadChat || !e.store.HasLiveWorker("gpu", "gpu", "model.chat") {
return provider.ChatResult{}, false, nil
}
wait := cfg.Worker.DistributedInferenceWaitS
if wait <= 0 {
wait = 180
}
j, err := e.store.EnqueueJobSpec(store.JobSpec{
Type: "model.chat", Payload: map[string]any{"instructions": instructions, "input": input, "model": model, "max_output": maxOutput, "json_mode": jsonMode},
Priority: 100, ResourceClass: "gpu", RequiredCapabilities: []string{"gpu", "model.chat"},
MaxAttempts: 2, TimeoutSeconds: wait,
})
if err != nil {
return provider.ChatResult{}, true, err
}
jobCtx, cancel := context.WithTimeout(ctx, time.Duration(wait)*time.Second)
defer cancel()
out, err := e.waitDistributedJob(jobCtx, j.ID)
if err != nil {
return provider.ChatResult{}, true, err
}
return provider.ChatResult{Text: out.Text, Provider: out.Provider, Model: out.Model, NodeID: out.NodeID, Usage: provider.Usage{InputTokens: out.Usage.InputTokens, CachedTokens: out.Usage.CachedTokens, OutputTokens: out.Usage.OutputTokens}}, true, nil
}
@@ -0,0 +1,117 @@
package brain
import (
"context"
"errors"
"time"
"neuroforge/internal/core"
)
// PlanGraphBackfill keeps the associative graph converging toward a durable,
// multi-neighbor state. It deliberately uses a bounded queue so a large KB
// import cannot create an unbounded O(N^2) burst of relink payloads.
func (e *Engine) PlanGraphBackfill() (int, error) {
cfg := e.store.Config()
if !e.store.IsOrchestratorLeader() {
return 0, nil
}
wc := cfg.Worker
if !wc.GraphBackfillEnabled {
return 0, nil
}
if wc.RequireWorkerForGraph && !e.store.HasLiveWorker("cpu", "cpu", "vector.relink") {
return 0, nil
}
maxQueued := wc.GraphBackfillMaxQueued
if maxQueued <= 0 {
maxQueued = 512
}
pending := e.store.PendingJobCount("vector.relink")
if pending >= maxQueued {
return 0, nil
}
batch := wc.GraphBackfillBatchSize
if batch <= 0 {
batch = 64
}
if room := maxQueued - pending; batch > room {
batch = room
}
retryAfter := time.Duration(wc.GraphRetryAfterMinutes) * time.Minute
if retryAfter <= 0 {
retryAfter = 6 * time.Hour
}
items := e.store.GraphBackfillCandidates(batch, wc.GraphBackfillMinDegree, retryAfter)
planned := 0
for i := range items {
if _, err := e.enqueueRelink(&items[i]); err != nil {
return planned, err
}
planned++
}
return planned, nil
}
// ApplyAndFinalizeJob is the durable second phase for jobs whose worker result
// mutates authoritative master state. The worker result is already in the WAL;
// apply can therefore be retried after a master crash without recomputing it.
func (e *Engine) ApplyAndFinalizeJob(j *core.Job) error {
if j == nil || j.Status != "apply_wait" {
return nil
}
err := e.ApplyJobResult(j)
if errors.Is(err, errObsoleteRelink) {
// Obsolete computation is not a failed mutation. Close the old job and
// leave the current memory unlinked so normal backfill schedules it again.
err = nil
}
_, finishErr := e.store.FinishMasterApply(j.ID, err)
if finishErr != nil {
return finishErr
}
return err
}
func (e *Engine) ApplyPendingJobResults(limit int) (applied, failed int) {
jobs := e.store.PendingMasterApplyJobs(limit)
for i := range jobs {
if err := e.ApplyAndFinalizeJob(&jobs[i]); err != nil {
failed++
} else {
applied++
}
}
return applied, failed
}
func (e *Engine) RunOrchestrator(ctx context.Context) {
var lastPrune time.Time
// Run quickly on boot so a previously imported knowledge corpus starts
// linking as soon as a capable worker has registered.
timer := time.NewTimer(2 * time.Second)
defer timer.Stop()
for {
select {
case <-ctx.Done():
return
case <-timer.C:
cfg := e.store.Config()
if !e.store.IsOrchestratorLeader() {
timer.Reset(5 * time.Second)
continue
}
_, _ = e.ApplyPendingJobResults(64)
if lastPrune.IsZero() || time.Since(lastPrune) >= 10*time.Minute {
_, _ = e.store.PruneTerminalJobs(time.Duration(cfg.Worker.JobRetentionHours)*time.Hour, cfg.Worker.MaxTerminalJobs)
lastPrune = time.Now().UTC()
}
interval := time.Duration(cfg.Worker.GraphBackfillIntervalS) * time.Second
if interval < 2*time.Second {
interval = 10 * time.Second
}
_, _ = e.PlanGraphBackfill()
timer.Reset(interval)
}
}
}
@@ -0,0 +1,200 @@
package brain
import (
"encoding/json"
"errors"
"testing"
"time"
"neuroforge/internal/core"
"neuroforge/internal/cost"
"neuroforge/internal/provider"
"neuroforge/internal/store"
)
func TestGraphBackfillPlansBoundedJobsAndDurablyAppliesResult(t *testing.T) {
s, err := store.New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
cfg := s.Config()
cfg.Worker.GraphBackfillEnabled = true
cfg.Worker.RequireWorkerForGraph = true
cfg.Worker.GraphBackfillBatchSize = 4
cfg.Worker.GraphBackfillMaxQueued = 8
cfg.Worker.GraphBackfillMinDegree = 2
cfg.Worker.GraphCandidateMultiplier = 4
cfg.Brain.RecallK = 3
cfg.Brain.MinSimilarity = .1
if err := s.UpdateConfig(cfg); err != nil {
t.Fatal(err)
}
for i := 0; i < 20; i++ {
m := core.Memory{ID: store.NewID("kb"), Kind: "knowledge.chunk", MemoryType: core.MemorySemantic, Text: "knowledge", Vector: []float32{1, float32(i+1) / 100}}
if err := s.AddMemory(&m); err != nil {
t.Fatal(err)
}
}
if _, err := s.RegisterWorker(store.WorkerHeartbeat{ID: "cpu-1", ResourceClass: "cpu", Capabilities: []string{"cpu", "vector.relink"}, MaxConcurrency: 2}); err != nil {
t.Fatal(err)
}
e := New(s, provider.NewRouter(s), cost.New(s))
planned, err := e.PlanGraphBackfill()
if err != nil {
t.Fatal(err)
}
if planned != 4 || s.PendingJobCount("vector.relink") != 4 {
t.Fatalf("planned=%d pending=%d", planned, s.PendingJobCount("vector.relink"))
}
jobs := s.JobsSnapshot(10, "", "vector.relink")
if len(jobs) != 4 {
t.Fatalf("jobs=%d", len(jobs))
}
full, ok := s.Job(jobs[0].ID)
if !ok || !full.RequiresMasterApply {
t.Fatalf("job=%+v", full)
}
var payload relinkPayload
if err := json.Unmarshal(full.Payload, &payload); err != nil {
t.Fatal(err)
}
if len(payload.Candidates) == 0 || len(payload.Candidates) > 256 {
t.Fatalf("candidate payload is not bounded: %d", len(payload.Candidates))
}
worker := store.WorkerHeartbeat{ID: "cpu-1", ResourceClass: "cpu", Capabilities: []string{"cpu", "vector.relink"}, MaxConcurrency: 2}
claimed, err := s.ClaimJobForWorker(worker, time.Minute)
if err != nil || claimed == nil {
t.Fatalf("claim=%+v err=%v", claimed, err)
}
var p relinkPayload
if err := json.Unmarshal(claimed.Payload, &p); err != nil {
t.Fatal(err)
}
if len(p.Candidates) < 2 {
t.Fatalf("need at least two candidates: %+v", p)
}
result, _ := json.Marshal(map[string]any{"target_id": p.TargetID, "neighbors": []map[string]any{{"id": p.Candidates[0].ID, "similarity": .95}, {"id": p.Candidates[1].ID, "similarity": .90}}})
completed, err := s.CompleteJobLease(claimed.ID, worker.ID, claimed.LeaseToken, result, "")
if err != nil {
t.Fatal(err)
}
if completed.Status != "apply_wait" {
t.Fatalf("worker result must be durable before master apply: %+v", completed)
}
if err := e.ApplyAndFinalizeJob(completed); err != nil {
t.Fatal(err)
}
final, _ := s.Job(completed.ID)
if final.Status != "done" || final.ApplyAttempts != 1 {
t.Fatalf("final=%+v", final)
}
if degree := s.GraphDegree(p.TargetID); degree < 2 {
t.Fatalf("target degree=%d; expected real multi-neighbor graph", degree)
}
m, ok := s.GetMemory(p.TargetID)
if !ok || m.GraphLinkedAt.IsZero() || m.GraphVersion != m.Version {
t.Fatalf("graph linkage checkpoint missing: %+v", m)
}
}
func TestMasterApplyReplayDoesNotInflateSemanticEdge(t *testing.T) {
s, err := store.New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
for _, id := range []string{"A", "B"} {
m := core.Memory{ID: id, Kind: "knowledge.chunk", MemoryType: core.MemorySemantic, Text: id, Vector: []float32{1, .1}}
if err := s.AddMemory(&m); err != nil {
t.Fatal(err)
}
}
e := New(s, provider.NewRouter(s), cost.New(s))
a, _ := s.GetMemory("A")
b, _ := s.GetMemory("B")
payload, _ := json.Marshal(relinkPayload{
TargetID: "A", TargetVersion: a.Version, TargetFingerprint: vectorFingerprint(a.Vector), Target: append([]float32(nil), a.Vector...),
Candidates: []relinkCandidate{{ID: "B", Version: b.Version, Fingerprint: vectorFingerprint(b.Vector), Vector: append([]float32(nil), b.Vector...)}},
K: 1, MinSimilarity: .1,
})
result, _ := json.Marshal(map[string]any{"target_id": "A", "neighbors": []map[string]any{{"id": "B", "similarity": .9}}})
j := &core.Job{Type: "vector.relink", Status: "apply_wait", Payload: payload, Result: result}
if err := e.ApplyJobResult(j); err != nil {
t.Fatal(err)
}
if err := e.ApplyJobResult(j); err != nil {
t.Fatal(err)
}
g := s.KnowledgeGraph("A", 1, 10)
if len(g.Edges) != 1 {
t.Fatalf("edges=%+v", g.Edges)
}
if g.Edges[0].Activations != 1 {
t.Fatalf("retry replay inflated activations: %+v", g.Edges[0])
}
}
func TestRelinkResultFromDeletedRecreatedMemoryIsObsolete(t *testing.T) {
s, err := store.New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
for _, m := range []core.Memory{
{ID: "A", Kind: "knowledge.chunk", MemoryType: core.MemorySemantic, Text: "old A", Vector: []float32{1, .1}},
{ID: "B", Kind: "knowledge.chunk", MemoryType: core.MemorySemantic, Text: "B", Vector: []float32{1, .2}},
} {
mm := m
if err := s.AddMemory(&mm); err != nil {
t.Fatal(err)
}
}
e := New(s, provider.NewRouter(s), cost.New(s))
job, err := e.enqueueRelink(mustMemory(t, s, "A"))
if err != nil {
t.Fatal(err)
}
var p relinkPayload
if err := json.Unmarshal(job.Payload, &p); err != nil {
t.Fatal(err)
}
if len(p.Candidates) == 0 {
t.Fatal("expected relink candidate")
}
// Simulate a knowledge sync replacing the stable ID while an old worker is
// still computing. Delete removes old edges; the re-created memory may reuse
// version 1, so the vector fingerprint is the required fencing signal.
if err := s.DeleteMemory("A"); err != nil {
t.Fatal(err)
}
recreated := core.Memory{ID: "A", Kind: "knowledge.chunk", MemoryType: core.MemorySemantic, Text: "new A", Vector: []float32{.1, 1}}
if err := s.AddMemory(&recreated); err != nil {
t.Fatal(err)
}
result, _ := json.Marshal(map[string]any{"target_id": "A", "neighbors": []map[string]any{{"id": p.Candidates[0].ID, "similarity": .99}}})
job.Status = "apply_wait"
job.Result = result
if err := e.ApplyJobResult(job); !errors.Is(err, errObsoleteRelink) {
t.Fatalf("expected obsolete relink result, got %v", err)
}
if got := s.GraphDegree("A"); got != 0 {
t.Fatalf("obsolete worker result mutated recreated memory graph: degree=%d", got)
}
m := mustMemory(t, s, "A")
if !m.GraphLinkedAt.IsZero() || m.GraphVersion != 0 {
t.Fatalf("recreated memory was incorrectly marked linked: %+v", m)
}
}
func mustMemory(t *testing.T, s *store.Store, id string) *core.Memory {
t.Helper()
m, ok := s.GetMemory(id)
if !ok {
t.Fatalf("memory %s missing", id)
}
return m
}
+97 -11
View File
@@ -173,6 +173,10 @@ type Config struct {
DecayPerDay float64 `json:"decay_per_day"`
MaxSynapseWeight float64 `json:"max_synapse_weight"`
GraphBonus float64 `json:"graph_bonus"`
GraphMaxHops int `json:"graph_max_hops"`
GraphHopDecay float64 `json:"graph_hop_decay"`
GraphMaxExpansion int `json:"graph_max_expansion"`
GraphMinEdgeWeight float64 `json:"graph_min_edge_weight"`
MaxContextMemories int `json:"max_context_memories"`
AutoLearn bool `json:"auto_learn"`
ExternalRelinkWorker bool `json:"external_relink_worker"`
@@ -388,7 +392,27 @@ type Config struct {
} `json:"api"`
Worker struct {
LeaseSeconds int `json:"lease_seconds"`
LeaseSeconds int `json:"lease_seconds"`
HeartbeatSeconds int `json:"heartbeat_seconds"`
StaleAfterSeconds int `json:"stale_after_seconds"`
DefaultMaxAttempts int `json:"default_max_attempts"`
RetryBackoffSeconds int `json:"retry_backoff_seconds"`
MaxQueuedJobs int `json:"max_queued_jobs"`
MasterApplyMaxAttempts int `json:"master_apply_max_attempts"`
MasterApplyBackoffSeconds int `json:"master_apply_backoff_seconds"`
JobRetentionHours int `json:"job_retention_hours"`
MaxTerminalJobs int `json:"max_terminal_jobs"`
GraphBackfillEnabled bool `json:"graph_backfill_enabled"`
GraphBackfillIntervalS int `json:"graph_backfill_interval_seconds"`
GraphBackfillBatchSize int `json:"graph_backfill_batch_size"`
GraphBackfillMaxQueued int `json:"graph_backfill_max_queued"`
GraphBackfillMinDegree int `json:"graph_backfill_min_degree"`
GraphCandidateMultiplier int `json:"graph_candidate_multiplier"`
GraphRetryAfterMinutes int `json:"graph_retry_after_minutes"`
RequireWorkerForGraph bool `json:"require_worker_for_graph"`
OffloadChat bool `json:"offload_chat"`
OffloadEmbeddings bool `json:"offload_embeddings"`
DistributedInferenceWaitS int `json:"distributed_inference_wait_seconds"`
} `json:"worker"`
}
@@ -469,6 +493,9 @@ type Memory struct {
ConsolidationCount int `json:"consolidation_count,omitempty"`
EvidenceSourceIDs []string `json:"evidence_source_ids,omitempty"`
EvidenceCount int `json:"evidence_count,omitempty"`
GraphLinkedAt time.Time `json:"graph_linked_at,omitempty"`
GraphDegree int `json:"graph_degree,omitempty"`
GraphVersion int64 `json:"graph_version,omitempty"`
Provenance MemoryProvenance `json:"provenance,omitempty"`
}
@@ -477,6 +504,7 @@ type Synapse struct {
B string `json:"b"`
Weight float64 `json:"weight"`
Similarity float64 `json:"similarity"`
Relations []string `json:"relations,omitempty"`
Activations int64 `json:"activations"`
LastUpdated time.Time `json:"last_updated"`
}
@@ -494,16 +522,50 @@ type UsageEvent struct {
}
type Job struct {
ID string `json:"id"`
Type string `json:"type"`
Payload json.RawMessage `json:"payload"`
Result json.RawMessage `json:"result,omitempty"`
Status string `json:"status"`
ClaimedBy string `json:"claimed_by,omitempty"`
LeaseUntil time.Time `json:"lease_until,omitempty"`
Error string `json:"error,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
ID string `json:"id"`
Type string `json:"type"`
Payload json.RawMessage `json:"payload"`
Result json.RawMessage `json:"result,omitempty"`
Status string `json:"status"`
Priority int `json:"priority,omitempty"`
ResourceClass string `json:"resource_class,omitempty"`
RequiredCapabilities []string `json:"required_capabilities,omitempty"`
IdempotencyKey string `json:"idempotency_key,omitempty"`
ParentJobID string `json:"parent_job_id,omitempty"`
DependsOn []string `json:"depends_on,omitempty"`
Attempts int `json:"attempts,omitempty"`
MaxAttempts int `json:"max_attempts,omitempty"`
BackoffSeconds int `json:"backoff_seconds,omitempty"`
TimeoutSeconds int `json:"timeout_seconds,omitempty"`
RequiresMasterApply bool `json:"requires_master_apply,omitempty"`
ApplyAttempts int `json:"apply_attempts,omitempty"`
MaxApplyAttempts int `json:"max_apply_attempts,omitempty"`
ApplyBackoffSeconds int `json:"apply_backoff_seconds,omitempty"`
ApplyNextAttemptAt time.Time `json:"apply_next_attempt_at,omitempty"`
ApplyError string `json:"apply_error,omitempty"`
ClaimedBy string `json:"claimed_by,omitempty"`
LeaseToken string `json:"lease_token,omitempty"`
LeaseUntil time.Time `json:"lease_until,omitempty"`
NextAttemptAt time.Time `json:"next_attempt_at,omitempty"`
StartedAt time.Time `json:"started_at,omitempty"`
FinishedAt time.Time `json:"finished_at,omitempty"`
Error string `json:"error,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}
type WorkerState struct {
ID string `json:"id"`
ResourceClass string `json:"resource_class"`
Capabilities []string `json:"capabilities"`
Labels map[string]string `json:"labels,omitempty"`
MaxConcurrency int `json:"max_concurrency"`
Inflight int `json:"inflight"`
Version string `json:"version,omitempty"`
Hostname string `json:"hostname,omitempty"`
LastHeartbeat time.Time `json:"last_heartbeat"`
RegisteredAt time.Time `json:"registered_at"`
Status string `json:"status"`
}
type MaintenanceStatus struct {
@@ -688,6 +750,10 @@ func DefaultConfig() Config {
c.Brain.DecayPerDay = 0.01
c.Brain.MaxSynapseWeight = 4.0
c.Brain.GraphBonus = 0.15
c.Brain.GraphMaxHops = 3
c.Brain.GraphHopDecay = 0.60
c.Brain.GraphMaxExpansion = 64
c.Brain.GraphMinEdgeWeight = 0.05
c.Brain.MaxContextMemories = 8
c.Brain.AutoLearn = true
c.Brain.ExternalRelinkWorker = true
@@ -850,5 +916,25 @@ func DefaultConfig() Config {
c.Cluster.LogSegmentBytes = 64 << 20
c.API.RequireKey = true
c.Worker.LeaseSeconds = 120
c.Worker.HeartbeatSeconds = 15
c.Worker.StaleAfterSeconds = 60
c.Worker.DefaultMaxAttempts = 3
c.Worker.RetryBackoffSeconds = 15
c.Worker.MaxQueuedJobs = 5000
c.Worker.MasterApplyMaxAttempts = 5
c.Worker.MasterApplyBackoffSeconds = 5
c.Worker.JobRetentionHours = 168
c.Worker.MaxTerminalJobs = 20000
c.Worker.GraphBackfillEnabled = true
c.Worker.GraphBackfillIntervalS = 10
c.Worker.GraphBackfillBatchSize = 64
c.Worker.GraphBackfillMaxQueued = 512
c.Worker.GraphBackfillMinDegree = 3
c.Worker.GraphCandidateMultiplier = 6
c.Worker.GraphRetryAfterMinutes = 360
c.Worker.RequireWorkerForGraph = true
c.Worker.OffloadChat = false
c.Worker.OffloadEmbeddings = false
c.Worker.DistributedInferenceWaitS = 180
return c
}
+104 -14
View File
@@ -92,6 +92,8 @@ func (s *Server) routes() {
s.mux.Handle("POST /api/v1/integrations/outcomes/search", s.integrationAuth(http.HandlerFunc(s.integrationValidatedOutcomeSearch)))
s.mux.Handle("GET /api/v1/integrations/graph/research", s.controlReadAuth(http.HandlerFunc(s.integrationResearchGraph)))
s.mux.Handle("GET /api/v1/integrations/graph/brain", s.controlReadAuth(http.HandlerFunc(s.integrationBrainGraph)))
s.mux.Handle("GET /api/v1/integrations/graph/status", s.controlReadAuth(http.HandlerFunc(s.integrationGraphStatus)))
s.mux.Handle("GET /api/v1/integrations/orchestrator/status", s.controlReadAuth(http.HandlerFunc(s.integrationOrchestratorStatus)))
s.mux.Handle("POST /internal/v1/cluster/request-vote", s.clusterAuth(http.HandlerFunc(s.clusterRequestVote)))
s.mux.Handle("POST /internal/v1/cluster/heartbeat", s.clusterAuth(http.HandlerFunc(s.clusterHeartbeat)))
@@ -102,6 +104,8 @@ func (s *Server) routes() {
s.mux.Handle("GET /internal/v1/cluster/decision/{id}", s.clusterAuth(http.HandlerFunc(s.clusterDecision)))
s.mux.Handle("GET /internal/v1/cluster/status", s.clusterAuth(http.HandlerFunc(s.clusterStatus)))
s.mux.Handle("POST /api/v1/worker/register", s.workerAuth(http.HandlerFunc(s.workerRegister)))
s.mux.Handle("POST /api/v1/worker/heartbeat", s.workerAuth(http.HandlerFunc(s.workerHeartbeat)))
s.mux.Handle("POST /api/v1/worker/claim", s.workerAuth(http.HandlerFunc(s.workerClaim)))
s.mux.Handle("POST /api/v1/worker/complete", s.workerAuth(http.HandlerFunc(s.workerComplete)))
@@ -138,6 +142,12 @@ func (s *Server) routes() {
s.mux.Handle("GET /admin/api/knowledge/memories", s.adminAuth(http.HandlerFunc(s.adminKnowledgeMemories)))
s.mux.Handle("GET /admin/api/knowledge/memory/{id}", s.adminAuth(http.HandlerFunc(s.adminKnowledgeMemory)))
s.mux.Handle("GET /admin/api/knowledge/graph", s.adminAuth(http.HandlerFunc(s.adminKnowledgeGraph)))
s.mux.Handle("GET /admin/api/graph/status", s.adminAuth(http.HandlerFunc(s.adminGraphStatus)))
s.mux.Handle("POST /admin/api/graph/backfill", s.adminAuth(http.HandlerFunc(s.adminGraphBackfill)))
s.mux.Handle("GET /admin/api/orchestrator/status", s.adminAuth(http.HandlerFunc(s.adminOrchestratorStatus)))
s.mux.Handle("GET /admin/api/orchestrator/jobs", s.adminAuth(http.HandlerFunc(s.adminOrchestratorJobs)))
s.mux.Handle("POST /admin/api/orchestrator/jobs/{id}/retry", s.adminAuth(http.HandlerFunc(s.adminOrchestratorRetry)))
s.mux.Handle("POST /admin/api/orchestrator/jobs/{id}/cancel", s.adminAuth(http.HandlerFunc(s.adminOrchestratorCancel)))
s.mux.Handle("GET /admin/api/knowledge/events", s.adminAuth(http.HandlerFunc(s.adminKnowledgeEvents)))
s.mux.Handle("POST /admin/api/knowledge/search", s.adminAuth(http.HandlerFunc(s.adminKnowledgeSearch)))
s.mux.Handle("GET /admin/api/learning-policy", s.adminAuth(http.HandlerFunc(s.adminGetLearningPolicy)))
@@ -418,10 +428,84 @@ func (s *Server) stats(w http.ResponseWriter, r *http.Request) {
s.json(w, 200, map[string]any{"stats": s.store.Stats(), "cost": s.cost.Totals()})
}
func (s *Server) workerClaim(w http.ResponseWriter, r *http.Request) {
var q struct {
WorkerID string `json:"worker_id"`
func (s *Server) integrationGraphStatus(w http.ResponseWriter, r *http.Request) {
cfg := s.store.Config()
s.json(w, http.StatusOK, map[string]any{
"stats": s.store.GraphStats(),
"config": map[string]any{
"max_hops": cfg.Brain.GraphMaxHops, "hop_decay": cfg.Brain.GraphHopDecay,
"max_expansion": cfg.Brain.GraphMaxExpansion, "min_edge_weight": cfg.Brain.GraphMinEdgeWeight,
"backfill_enabled": cfg.Worker.GraphBackfillEnabled, "backfill_min_degree": cfg.Worker.GraphBackfillMinDegree,
},
})
}
func (s *Server) integrationOrchestratorStatus(w http.ResponseWriter, r *http.Request) {
cfg := s.store.Config()
s.json(w, http.StatusOK, map[string]any{
"leader": s.store.IsOrchestratorLeader(),
"status": s.store.OrchestratorStatus(),
"config": map[string]any{
"lease_seconds": cfg.Worker.LeaseSeconds, "heartbeat_seconds": cfg.Worker.HeartbeatSeconds,
"stale_after_seconds": cfg.Worker.StaleAfterSeconds, "max_queued_jobs": cfg.Worker.MaxQueuedJobs,
"offload_chat": cfg.Worker.OffloadChat, "offload_embeddings": cfg.Worker.OffloadEmbeddings,
},
})
}
type workerRequest struct {
WorkerID string `json:"worker_id"`
ResourceClass string `json:"resource_class,omitempty"`
Capabilities []string `json:"capabilities,omitempty"`
Labels map[string]string `json:"labels,omitempty"`
MaxConcurrency int `json:"max_concurrency,omitempty"`
Version string `json:"version,omitempty"`
Hostname string `json:"hostname,omitempty"`
ActiveLeases map[string]string `json:"active_leases,omitempty"`
}
func workerHeartbeatFromRequest(q workerRequest) store.WorkerHeartbeat {
return store.WorkerHeartbeat{
ID: q.WorkerID, ResourceClass: q.ResourceClass, Capabilities: q.Capabilities,
Labels: q.Labels, MaxConcurrency: q.MaxConcurrency, Version: q.Version,
Hostname: q.Hostname, ActiveLeases: q.ActiveLeases,
}
}
func (s *Server) workerRegister(w http.ResponseWriter, r *http.Request) {
var q workerRequest
if err := decode(r, &q); err != nil {
s.err(w, http.StatusBadRequest, err)
return
}
state, err := s.store.RegisterWorker(workerHeartbeatFromRequest(q))
if err != nil {
s.err(w, http.StatusBadRequest, err)
return
}
s.json(w, http.StatusOK, state)
}
func (s *Server) workerHeartbeat(w http.ResponseWriter, r *http.Request) {
var q workerRequest
if err := decode(r, &q); err != nil {
s.err(w, http.StatusBadRequest, err)
return
}
lease := s.store.Config().Worker.LeaseSeconds
if lease < 10 {
lease = 120
}
state, err := s.store.HeartbeatWorker(workerHeartbeatFromRequest(q), time.Duration(lease)*time.Second)
if err != nil {
s.err(w, http.StatusBadRequest, err)
return
}
s.json(w, http.StatusOK, state)
}
func (s *Server) workerClaim(w http.ResponseWriter, r *http.Request) {
var q workerRequest
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
@@ -434,7 +518,7 @@ func (s *Server) workerClaim(w http.ResponseWriter, r *http.Request) {
if lease < 10 {
lease = 120
}
j, err := s.store.ClaimJob(q.WorkerID, time.Duration(lease)*time.Second)
j, err := s.store.ClaimJobForWorker(workerHeartbeatFromRequest(q), time.Duration(lease)*time.Second)
if err != nil {
s.err(w, 500, err)
return
@@ -445,27 +529,33 @@ func (s *Server) workerClaim(w http.ResponseWriter, r *http.Request) {
}
s.json(w, 200, j)
}
func (s *Server) workerComplete(w http.ResponseWriter, r *http.Request) {
var q struct {
WorkerID string `json:"worker_id"`
JobID string `json:"job_id"`
Result json.RawMessage `json:"result"`
Error string `json:"error"`
WorkerID string `json:"worker_id"`
JobID string `json:"job_id"`
LeaseToken string `json:"lease_token,omitempty"`
Result json.RawMessage `json:"result"`
Error string `json:"error"`
}
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
j, err := s.store.CompleteJob(q.JobID, q.WorkerID, q.Result, q.Error)
j, err := s.store.CompleteJobLease(q.JobID, q.WorkerID, q.LeaseToken, q.Result, q.Error)
if err != nil {
s.err(w, 400, err)
s.err(w, 409, err)
return
}
if err := s.brain.ApplyJobResult(j); err != nil {
s.err(w, 500, err)
return
if j.Status == "apply_wait" && s.store.IsOrchestratorLeader() {
// Best-effort eager second phase. Any failure remains durable as apply_wait
// and is retried by the orchestrator; the worker must not recompute it.
_ = s.brain.ApplyAndFinalizeJob(j)
if cur, ok := s.store.Job(j.ID); ok {
j = cur
}
}
s.json(w, 200, map[string]bool{"ok": true})
s.json(w, 200, map[string]any{"ok": true, "status": j.Status, "attempts": j.Attempts, "apply_attempts": j.ApplyAttempts, "apply_error": j.ApplyError, "next_attempt_at": j.NextAttemptAt, "apply_next_attempt_at": j.ApplyNextAttemptAt})
}
func (s *Server) adminStatus(w http.ResponseWriter, r *http.Request) {
@@ -115,3 +115,69 @@ func (s *Server) adminPutLearningPolicy(w http.ResponseWriter, r *http.Request)
})
s.json(w, http.StatusOK, q)
}
func (s *Server) adminGraphStatus(w http.ResponseWriter, r *http.Request) {
s.json(w, http.StatusOK, map[string]any{
"graph": s.store.GraphStats(),
"orchestrator": s.store.OrchestratorStatus(),
"config": map[string]any{
"max_hops": s.store.Config().Brain.GraphMaxHops,
"hop_decay": s.store.Config().Brain.GraphHopDecay,
"max_expansion": s.store.Config().Brain.GraphMaxExpansion,
"min_edge_weight": s.store.Config().Brain.GraphMinEdgeWeight,
"backfill_enabled": s.store.Config().Worker.GraphBackfillEnabled,
"backfill_min_degree": s.store.Config().Worker.GraphBackfillMinDegree,
"backfill_max_queued": s.store.Config().Worker.GraphBackfillMaxQueued,
},
})
}
func (s *Server) adminGraphBackfill(w http.ResponseWriter, r *http.Request) {
planned, err := s.brain.PlanGraphBackfill()
if err != nil {
s.err(w, http.StatusInternalServerError, err)
return
}
s.json(w, http.StatusAccepted, map[string]any{"ok": true, "planned": planned, "graph": s.store.GraphStats()})
}
func (s *Server) adminOrchestratorStatus(w http.ResponseWriter, r *http.Request) {
s.json(w, http.StatusOK, s.store.OrchestratorStatus())
}
func (s *Server) adminOrchestratorJobs(w http.ResponseWriter, r *http.Request) {
limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
s.json(w, http.StatusOK, s.store.JobsSnapshot(limit, r.URL.Query().Get("status"), r.URL.Query().Get("type")))
}
func (s *Server) adminOrchestratorRetry(w http.ResponseWriter, r *http.Request) {
job, err := s.store.RetryJob(r.PathValue("id"))
if err != nil {
s.err(w, http.StatusConflict, err)
return
}
s.json(w, http.StatusOK, job)
}
func (s *Server) adminOrchestratorCancel(w http.ResponseWriter, r *http.Request) {
var q struct {
Reason string `json:"reason"`
}
// Empty bodies are valid for an operator cancellation; malformed non-empty
// JSON is not.
if r.ContentLength != 0 {
if err := decode(r, &q); err != nil {
s.err(w, http.StatusBadRequest, err)
return
}
}
if strings.TrimSpace(q.Reason) == "" {
q.Reason = "canceled by administrator"
}
if err := s.store.CancelJob(r.PathValue("id"), q.Reason); err != nil {
s.err(w, http.StatusNotFound, err)
return
}
job, _ := s.store.Job(r.PathValue("id"))
s.json(w, http.StatusOK, job)
}
@@ -242,6 +242,8 @@ func (s *Server) metricsEndpoint(w http.ResponseWriter, r *http.Request) {
}
st := s.store.ObservabilitySnapshot()
orch := s.store.OrchestratorStatus()
graph := s.store.GraphStats()
costs := s.cost.Totals()
cfg := s.store.Config()
rt := currentRuntimeSnapshot()
@@ -269,11 +271,43 @@ func (s *Server) metricsEndpoint(w http.ResponseWriter, r *http.Request) {
promHeader(&b, "neuroforge_knowledge_events", "Current number of retained explainability/knowledge events.", "gauge")
promSample(&b, "neuroforge_knowledge_events", st.KnowledgeEvents)
promHeader(&b, "neuroforge_jobs", "Current worker jobs by bounded status class.", "gauge")
promSample(&b, "neuroforge_jobs", st.JobsQueued, "status", "queued")
promSample(&b, "neuroforge_jobs", st.JobsClaimed, "status", "claimed")
promSample(&b, "neuroforge_jobs", st.JobsDone, "status", "done")
promSample(&b, "neuroforge_jobs", st.JobsFailed, "status", "failed")
promHeader(&b, "neuroforge_jobs", "Current durable orchestrator jobs by status.", "gauge")
statuses := []string{"queued", "claimed", "retry_wait", "blocked", "apply_wait", "done", "failed", "canceled"}
for _, status := range statuses {
promSample(&b, "neuroforge_jobs", orch.JobsByStatus[status], "status", status)
}
promHeader(&b, "neuroforge_jobs_by_resource", "Current durable orchestrator jobs by resource class.", "gauge")
for resource, n := range orch.JobsByResource {
promSample(&b, "neuroforge_jobs_by_resource", n, "resource", resource)
}
promHeader(&b, "neuroforge_workers", "Registered subagents by resource class and liveness state.", "gauge")
workerCounts := map[string]int{}
for _, worker := range orch.Workers {
workerCounts[worker.ResourceClass+"\x00"+worker.Status]++
}
for key, n := range workerCounts {
parts := strings.SplitN(key, "\x00", 2)
promSample(&b, "neuroforge_workers", n, "resource", parts[0], "status", parts[1])
}
promHeader(&b, "neuroforge_worker_inflight", "Current claimed jobs per registered subagent.", "gauge")
promHeader(&b, "neuroforge_worker_capacity", "Configured concurrency capacity per registered subagent.", "gauge")
for _, worker := range orch.Workers {
promSample(&b, "neuroforge_worker_inflight", worker.Inflight, "worker", worker.ID, "resource", worker.ResourceClass)
promSample(&b, "neuroforge_worker_capacity", worker.MaxConcurrency, "worker", worker.ID, "resource", worker.ResourceClass)
}
promHeader(&b, "neuroforge_graph_memories", "Knowledge memories by graph linkage class.", "gauge")
promSample(&b, "neuroforge_graph_memories", graph.LinkedMemories, "state", "linked")
promSample(&b, "neuroforge_graph_memories", graph.IsolatedMemories, "state", "isolated")
promSample(&b, "neuroforge_graph_memories", graph.MultiLinked, "state", "multi_linked")
promHeader(&b, "neuroforge_graph_max_degree", "Maximum associative graph degree.", "gauge")
promSample(&b, "neuroforge_graph_max_degree", graph.MaxDegree)
promHeader(&b, "neuroforge_graph_average_degree", "Average associative graph degree.", "gauge")
promSample(&b, "neuroforge_graph_average_degree", strconv.FormatFloat(graph.AverageDegree, 'f', 6, 64))
promHeader(&b, "neuroforge_graph_connected_components", "Number of connected components in the associative graph.", "gauge")
promSample(&b, "neuroforge_graph_connected_components", graph.ConnectedComponents)
promHeader(&b, "neuroforge_graph_largest_component", "Memories in the largest connected component.", "gauge")
promSample(&b, "neuroforge_graph_largest_component", graph.LargestComponent)
promHeader(&b, "neuroforge_hnsw_nodes", "Current number of vectors in hot HNSW indexes.", "gauge")
promSample(&b, "neuroforge_hnsw_nodes", st.HNSWNodes)
@@ -129,6 +129,7 @@ func (s *Store) DeleteMemoriesBatch(ids []string) error {
}
for key, syn := range s.state.Synapses {
if syn.A == id || syn.B == id {
s.unindexSynapseLocked(syn)
delete(s.state.Synapses, key)
}
}
@@ -0,0 +1,405 @@
package store
import (
"sort"
"strings"
"time"
"neuroforge/internal/core"
"neuroforge/internal/vector"
)
// GraphStats exposes structural graph health, not just raw edge counts. A real
// associative graph is expected to have nodes with degree >1 and multi-hop
// reachability; these metrics make that observable.
type GraphStats struct {
Memories int `json:"memories"`
Synapses int `json:"synapses"`
LinkedMemories int `json:"linked_memories"`
IsolatedMemories int `json:"isolated_memories"`
MultiLinked int `json:"multi_linked_memories"`
MaxDegree int `json:"max_degree"`
AverageDegree float64 `json:"average_degree"`
ConnectedComponents int `json:"connected_components"`
LargestComponent int `json:"largest_component"`
}
func (s *Store) rebuildSynapseAdjLocked() {
s.synapseAdj = make(map[string]map[string]*core.Synapse)
for _, syn := range s.state.Synapses {
if syn == nil || syn.A == "" || syn.B == "" || syn.A == syn.B {
continue
}
s.indexSynapseLocked(syn)
}
}
func (s *Store) indexSynapseLocked(syn *core.Synapse) {
if syn == nil || syn.A == "" || syn.B == "" || syn.A == syn.B {
return
}
if s.synapseAdj == nil {
s.synapseAdj = make(map[string]map[string]*core.Synapse)
}
if s.synapseAdj[syn.A] == nil {
s.synapseAdj[syn.A] = make(map[string]*core.Synapse)
}
if s.synapseAdj[syn.B] == nil {
s.synapseAdj[syn.B] = make(map[string]*core.Synapse)
}
s.synapseAdj[syn.A][syn.B] = syn
s.synapseAdj[syn.B][syn.A] = syn
}
func (s *Store) unindexSynapseLocked(syn *core.Synapse) {
if syn == nil || s.synapseAdj == nil {
return
}
if m := s.synapseAdj[syn.A]; m != nil {
delete(m, syn.B)
if len(m) == 0 {
delete(s.synapseAdj, syn.A)
}
}
if m := s.synapseAdj[syn.B]; m != nil {
delete(m, syn.A)
if len(m) == 0 {
delete(s.synapseAdj, syn.B)
}
}
}
func appendUniqueRelation(xs []string, value string) []string {
for _, x := range xs {
if x == value {
return xs
}
}
return append(xs, value)
}
// ReinforceRelation adds or reinforces an associative edge while preserving the
// reason(s) the edge exists. Multiple relations can coexist on the same pair.
func (s *Store) ReinforceRelation(a, b, relation string, similarity, delta, decayPerDay, maxWeight float64) error {
if a == "" || b == "" || a == b {
return nil
}
s.mu.Lock()
defer s.mu.Unlock()
if s.state.Memories[a] == nil || s.state.Memories[b] == nil {
return nil
}
key := edgeKey(a, b)
now := time.Now().UTC()
syn, ok := s.state.Synapses[key]
if !ok {
syn = &core.Synapse{A: a, B: b, Similarity: similarity, LastUpdated: now}
s.state.Synapses[key] = syn
s.indexSynapseLocked(syn)
}
if relation == "" {
relation = "association"
}
syn.Relations = appendUniqueRelation(syn.Relations, relation)
days := now.Sub(syn.LastUpdated).Hours() / 24
if days > 0 && decayPerDay > 0 {
syn.Weight *= pow(1-decayPerDay, days)
}
syn.Weight = vector.Clamp(syn.Weight+delta, -maxWeight, maxWeight)
if similarity > syn.Similarity {
syn.Similarity = similarity
}
syn.Activations++
syn.LastUpdated = now
return s.commitLocked("synapse.upsert", *syn)
}
func (s *Store) GraphDegree(id string) int {
s.mu.RLock()
defer s.mu.RUnlock()
return len(s.synapseAdj[id])
}
func (s *Store) GraphStats() GraphStats {
s.mu.RLock()
defer s.mu.RUnlock()
out := GraphStats{Memories: len(s.state.Memories), Synapses: len(s.state.Synapses)}
if out.Memories == 0 {
return out
}
for id := range s.state.Memories {
degree := len(s.synapseAdj[id])
if degree == 0 {
out.IsolatedMemories++
} else {
out.LinkedMemories++
}
if degree > 1 {
out.MultiLinked++
}
if degree > out.MaxDegree {
out.MaxDegree = degree
}
}
out.AverageDegree = float64(out.Synapses*2) / float64(out.Memories)
visited := make(map[string]bool, out.Memories)
for id := range s.state.Memories {
if visited[id] {
continue
}
out.ConnectedComponents++
queue := []string{id}
visited[id] = true
size := 0
for len(queue) > 0 {
cur := queue[0]
queue = queue[1:]
size++
for next := range s.synapseAdj[cur] {
if !visited[next] && s.state.Memories[next] != nil {
visited[next] = true
queue = append(queue, next)
}
}
}
if size > out.LargestComponent {
out.LargestComponent = size
}
}
return out
}
// graphExpandLocked applies bounded multi-hop associative propagation to a
// vector result set. The graph remains pairwise at the edge level (as graphs
// naturally are) but retrieval can now follow chains across several memories.
func (s *Store) graphExpandLocked(q []float32, hits []SearchHit, k int, min, graphBonus float64) []SearchHit {
cfg := s.state.Config.Brain
if graphBonus <= 0 || len(hits) == 0 || len(s.state.Synapses) == 0 {
return hits
}
maxHops := cfg.GraphMaxHops
if maxHops <= 0 {
maxHops = 1
}
if maxHops > 6 {
maxHops = 6
}
decay := cfg.GraphHopDecay
if decay <= 0 || decay > 1 {
decay = 0.60
}
maxExpansion := cfg.GraphMaxExpansion
if maxExpansion <= 0 {
maxExpansion = 64
}
minEdge := cfg.GraphMinEdgeWeight
type frontierItem struct {
id string
boost float64
hop int
}
bestBoost := map[string]float64{}
queue := make([]frontierItem, 0, len(hits)*2)
for _, h := range hits {
bestBoost[h.Memory.ID] = 0
queue = append(queue, frontierItem{id: h.Memory.ID, boost: graphBonus, hop: 0})
}
expanded := 0
for len(queue) > 0 && expanded < maxExpansion {
cur := queue[0]
queue = queue[1:]
if cur.hop >= maxHops {
continue
}
neighbors := s.synapseAdj[cur.id]
if len(neighbors) == 0 {
continue
}
type edgeItem struct {
id string
syn *core.Synapse
}
edges := make([]edgeItem, 0, len(neighbors))
for id, syn := range neighbors {
if syn != nil && syn.Weight >= minEdge {
edges = append(edges, edgeItem{id: id, syn: syn})
}
}
sort.Slice(edges, func(i, j int) bool { return edges[i].syn.Weight > edges[j].syn.Weight })
for _, edge := range edges {
if expanded >= maxExpansion {
break
}
meta := s.state.Memories[edge.id]
if meta == nil || !memorySearchable(meta) {
continue
}
boost := cur.boost * edge.syn.Weight
if cur.hop > 0 {
boost *= decay
}
if old, ok := bestBoost[edge.id]; ok && old >= boost {
continue
}
bestBoost[edge.id] = boost
queue = append(queue, frontierItem{id: edge.id, boost: boost, hop: cur.hop + 1})
expanded++
}
}
byID := make(map[string]int, len(hits))
for i := range hits {
byID[hits[i].Memory.ID] = i
}
for id, boost := range bestBoost {
if boost <= 0 {
continue
}
if i, ok := byID[id]; ok {
hits[i].GraphBoost += boost
hits[i].Score += boost
continue
}
m, ok := s.fullMemoryForReadLocked(id)
if !ok || len(m.Vector) != len(q) {
continue
}
sim := vector.Cosine(q, m.Vector)
if sim < min {
continue
}
baseScore := sim * 0.5
hits = append(hits, SearchHit{Memory: cloneMemory(m), Similarity: sim, BaseScore: baseScore, GraphBoost: boost, Score: baseScore + boost, TypeWeight: 1, SalienceFactor: 1, ConfidenceFactor: 1, CandidateSource: "synapse-multihop"})
}
sort.Slice(hits, func(i, j int) bool { return hits[i].Score > hits[j].Score })
if len(hits) > k {
hits = hits[:k]
}
return hits
}
// GraphBackfillCandidates returns a bounded set of memories whose associative
// links are missing, stale for the current memory version, or remain below the
// configured minimum degree after the retry cooldown.
func (s *Store) GraphBackfillCandidates(limit, minDegree int, retryAfter time.Duration) []core.Memory {
if limit <= 0 {
return nil
}
if minDegree < 1 {
minDegree = 1
}
now := time.Now().UTC()
s.mu.RLock()
defer s.mu.RUnlock()
out := make([]core.Memory, 0, limit)
// Do not repeatedly return a target that already has a durable relink job
// in flight. Without this guard, the lexicographically first memories can
// occupy every planning batch until their first job finishes and starve the
// rest of a large imported corpus.
pendingTargets := make(map[string]bool)
for _, j := range s.state.Jobs {
if j == nil || j.Type != "vector.relink" {
continue
}
switch j.Status {
case "queued", "claimed", "retry_wait", "blocked", "apply_wait":
default:
continue
}
const prefix = "vector.relink:"
if !strings.HasPrefix(j.IdempotencyKey, prefix) {
continue
}
rest := strings.TrimPrefix(j.IdempotencyKey, prefix)
if cut := strings.LastIndex(rest, ":v"); cut > 0 {
pendingTargets[rest[:cut]] = true
}
}
ids := make([]string, 0, len(s.state.Memories))
for id := range s.state.Memories {
ids = append(ids, id)
}
sort.Strings(ids)
for _, id := range ids {
meta := s.state.Memories[id]
if pendingTargets[id] || meta == nil || !memorySearchable(meta) {
continue
}
m, ok := s.fullMemoryForReadLocked(id)
if !ok || len(m.Vector) == 0 {
continue
}
degree := len(s.synapseAdj[id])
needs := m.GraphVersion < m.Version || m.GraphLinkedAt.IsZero()
if !needs && degree < minDegree && (retryAfter <= 0 || now.Sub(m.GraphLinkedAt) >= retryAfter) {
needs = true
}
if !needs {
continue
}
m.GraphDegree = degree
out = append(out, cloneMemory(m))
if len(out) >= limit {
break
}
}
return out
}
func (s *Store) MarkGraphLinked(id string) error {
s.mu.Lock()
defer s.mu.Unlock()
m, ok := s.materializeMemoryLocked(id)
if !ok || m == nil {
return nil
}
m.GraphLinkedAt = time.Now().UTC()
m.GraphDegree = len(s.synapseAdj[id])
m.GraphVersion = m.Version
return s.commitLocked("memory.upsert", []core.Memory{cloneMemory(*m)})
}
// UpsertRelationEvidence records deterministic relation evidence without
// additive reinforcement. It is used for retryable master-apply phases where
// the same worker result may be replayed after a crash; replaying it must not
// inflate graph weights or activation counters.
func (s *Store) UpsertRelationEvidence(a, b, relation string, similarity, weight, maxWeight float64) error {
if a == "" || b == "" || a == b {
return nil
}
s.mu.Lock()
defer s.mu.Unlock()
if s.state.Memories[a] == nil || s.state.Memories[b] == nil {
return nil
}
if relation == "" {
relation = "association"
}
key := edgeKey(a, b)
now := time.Now().UTC()
syn, ok := s.state.Synapses[key]
if !ok {
syn = &core.Synapse{A: a, B: b, LastUpdated: now}
s.state.Synapses[key] = syn
s.indexSynapseLocked(syn)
}
beforeRelation := len(syn.Relations)
syn.Relations = appendUniqueRelation(syn.Relations, relation)
if similarity > syn.Similarity {
syn.Similarity = similarity
}
if maxWeight <= 0 {
maxWeight = 1
}
weight = vector.Clamp(weight, -maxWeight, maxWeight)
if weight > syn.Weight {
syn.Weight = weight
}
// Count creation of relation evidence, not crash/retry replays.
if len(syn.Relations) > beforeRelation || syn.Activations == 0 {
syn.Activations++
}
syn.LastUpdated = now
return s.commitLocked("synapse.upsert", *syn)
}
@@ -0,0 +1,103 @@
package store
import (
"testing"
"time"
"neuroforge/internal/core"
)
func TestGraphIsManyToManyAndTraversesMultipleHops(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
cfg := s.Config()
cfg.Brain.GraphMaxHops = 3
cfg.Brain.GraphHopDecay = .8
cfg.Brain.GraphMaxExpansion = 32
cfg.Brain.GraphMinEdgeWeight = .01
if err := s.UpdateConfig(cfg); err != nil {
t.Fatal(err)
}
items := []core.Memory{
{ID: "A", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "A", Vector: []float32{1, .01}},
{ID: "B", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "B", Vector: []float32{.99, .1}},
{ID: "C", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "C", Vector: []float32{.9, .3}},
{ID: "D", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "D", Vector: []float32{.8, .4}},
}
if err := s.AddMemoriesBatch(items); err != nil {
t.Fatal(err)
}
for _, e := range [][2]string{{"A", "B"}, {"B", "C"}, {"C", "D"}, {"B", "D"}} {
if err := s.ReinforceRelation(e[0], e[1], "semantic_similarity", .9, .8, 0, 1); err != nil {
t.Fatal(err)
}
}
// A second relation on the same pair proves relation metadata is not limited
// to a single 1:1 reason.
if err := s.ReinforceRelation("B", "C", "coactivation", .9, .05, 0, 1); err != nil {
t.Fatal(err)
}
st := s.GraphStats()
if st.Synapses != 4 || st.MultiLinked < 2 || st.MaxDegree < 3 || st.ConnectedComponents != 1 || st.LargestComponent != 4 {
t.Fatalf("graph stats do not prove an n:m connected graph: %+v", st)
}
s.mu.RLock()
seed, ok := s.fullMemoryForReadLocked("A")
if !ok {
s.mu.RUnlock()
t.Fatal("missing A")
}
hits := s.graphExpandLocked([]float32{1, .01}, []SearchHit{{Memory: cloneMemory(seed), Similarity: 1, Score: 1, BaseScore: 1}}, 4, .1, .25)
s.mu.RUnlock()
foundD := false
for _, h := range hits {
if h.Memory.ID == "D" && h.CandidateSource == "synapse-multihop" {
foundD = true
}
}
if !foundD {
t.Fatalf("multi-hop traversal did not surface D: %+v", hits)
}
g := s.KnowledgeGraph("B", 2, 32)
relationFound := false
for _, e := range g.Edges {
if (e.A == "B" && e.B == "C") || (e.A == "C" && e.B == "B") {
if len(e.Relations) >= 2 {
relationFound = true
}
}
}
if !relationFound {
t.Fatalf("typed/multi-relation edge not exposed: %+v", g.Edges)
}
}
func TestGraphBackfillSkipsTargetsAlreadyQueued(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
items := []core.Memory{
{ID: "A", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "A", Vector: []float32{1, 0}},
{ID: "B", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "B", Vector: []float32{.9, .1}},
{ID: "C", Kind: "knowledge", MemoryType: core.MemorySemantic, Text: "C", Vector: []float32{.8, .2}},
}
if err := s.AddMemoriesBatch(items); err != nil {
t.Fatal(err)
}
if _, err := s.EnqueueJobSpec(JobSpec{Type: "vector.relink", Payload: map[string]any{"target_id": "A"}, IdempotencyKey: "vector.relink:A:v1", ResourceClass: "cpu"}); err != nil {
t.Fatal(err)
}
got := s.GraphBackfillCandidates(2, 3, time.Hour)
if len(got) != 2 || got[0].ID == "A" || got[1].ID == "A" {
t.Fatalf("queued target was not skipped: %+v", got)
}
}
@@ -61,6 +61,7 @@ type KnowledgeList struct {
type KnowledgeEdge struct {
A string `json:"a"`
B string `json:"b"`
Relations []string `json:"relations,omitempty"`
Weight float64 `json:"weight"`
Similarity float64 `json:"similarity"`
Activations int64 `json:"activations"`
@@ -368,7 +369,7 @@ func (s *Store) KnowledgeGraph(center string, depth, maxNodes int) KnowledgeGrap
}
for _, syn := range s.state.Synapses {
if selected[syn.A] && selected[syn.B] {
out.Edges = append(out.Edges, KnowledgeEdge{A: syn.A, B: syn.B, Weight: syn.Weight, Similarity: syn.Similarity, Activations: syn.Activations, LastUpdated: syn.LastUpdated})
out.Edges = append(out.Edges, KnowledgeEdge{A: syn.A, B: syn.B, Relations: append([]string(nil), syn.Relations...), Weight: syn.Weight, Similarity: syn.Similarity, Activations: syn.Activations, LastUpdated: syn.LastUpdated})
}
}
sort.Slice(out.Edges, func(i, j int) bool { return out.Edges[i].Weight > out.Edges[j].Weight })
@@ -397,7 +398,7 @@ func (s *Store) KnowledgeMemoryDetail(id string) (KnowledgeMemoryDetail, bool) {
} else {
continue
}
out.Edges = append(out.Edges, KnowledgeEdge{A: syn.A, B: syn.B, Weight: syn.Weight, Similarity: syn.Similarity, Activations: syn.Activations, LastUpdated: syn.LastUpdated})
out.Edges = append(out.Edges, KnowledgeEdge{A: syn.A, B: syn.B, Relations: append([]string(nil), syn.Relations...), Weight: syn.Weight, Similarity: syn.Similarity, Activations: syn.Activations, LastUpdated: syn.LastUpdated})
if !seenNeighbor[other] {
if om, ok := s.fullMemoryForReadLocked(other); ok {
out.Neighbors = append(out.Neighbors, memoryPreview(om))
@@ -0,0 +1,841 @@
package store
import (
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
"time"
"neuroforge/internal/core"
)
// JobSpec describes a durable unit of work. Jobs are persisted through the
// normal NeuroForge WAL/checkpoint path so a master restart does not lose
// queued, leased, retrying, or dependency-blocked work.
type JobSpec struct {
Type string
Payload any
Priority int
ResourceClass string
RequiredCapabilities []string
IdempotencyKey string
ParentJobID string
DependsOn []string
MaxAttempts int
BackoffSeconds int
TimeoutSeconds int
RequiresMasterApply bool
MaxApplyAttempts int
ApplyBackoffSeconds int
}
type WorkerHeartbeat struct {
ID string
ResourceClass string
Capabilities []string
Labels map[string]string
MaxConcurrency int
Version string
Hostname string
ActiveLeases map[string]string // job_id -> lease_token
}
type OrchestratorStatus struct {
Workers []core.WorkerState `json:"workers"`
JobsByStatus map[string]int `json:"jobs_by_status"`
JobsByType map[string]int `json:"jobs_by_type"`
JobsByResource map[string]int `json:"jobs_by_resource"`
OldestQueued time.Time `json:"oldest_queued,omitempty"`
Queued int `json:"queued"`
Claimed int `json:"claimed"`
Retrying int `json:"retrying"`
Applying int `json:"applying"`
Blocked int `json:"blocked"`
Failed int `json:"failed"`
Done int `json:"done"`
}
func normalizeCapabilities(in []string) []string {
seen := map[string]bool{}
out := make([]string, 0, len(in))
for _, x := range in {
x = strings.ToLower(strings.TrimSpace(x))
if x == "" || seen[x] {
continue
}
seen[x] = true
out = append(out, x)
}
sort.Strings(out)
return out
}
func capabilitiesContain(have []string, need []string) bool {
if len(need) == 0 {
return true
}
m := make(map[string]bool, len(have))
for _, x := range have {
m[strings.ToLower(strings.TrimSpace(x))] = true
}
for _, x := range need {
if !m[strings.ToLower(strings.TrimSpace(x))] {
return false
}
}
return true
}
func (s *Store) EnqueueJob(kind string, payload any) (*core.Job, error) {
return s.EnqueueJobSpec(JobSpec{Type: kind, Payload: payload})
}
func (s *Store) EnqueueJobSpec(spec JobSpec) (*core.Job, error) {
spec.Type = strings.TrimSpace(spec.Type)
if spec.Type == "" {
return nil, errors.New("job type is required")
}
b, err := json.Marshal(spec.Payload)
if err != nil {
return nil, err
}
s.mu.Lock()
defer s.mu.Unlock()
cfg := s.state.Config.Worker
if cfg.MaxQueuedJobs > 0 {
pending := 0
for _, j := range s.state.Jobs {
if j != nil && (j.Status == "queued" || j.Status == "claimed" || j.Status == "retry_wait" || j.Status == "blocked" || j.Status == "apply_wait") {
pending++
}
}
if pending >= cfg.MaxQueuedJobs {
return nil, fmt.Errorf("orchestrator queue full: %d >= %d", pending, cfg.MaxQueuedJobs)
}
}
if key := strings.TrimSpace(spec.IdempotencyKey); key != "" {
for _, j := range s.state.Jobs {
if j != nil && j.IdempotencyKey == key && (j.Status == "queued" || j.Status == "claimed" || j.Status == "retry_wait" || j.Status == "blocked" || j.Status == "apply_wait") {
cp := *j
return &cp, nil
}
}
}
// Dependencies must already exist. Accepting dangling IDs would create jobs
// that can remain blocked forever and makes convergence impossible to reason
// about after operator typos or partial planning failures.
seenDeps := make(map[string]struct{}, len(spec.DependsOn))
cleanDeps := make([]string, 0, len(spec.DependsOn))
for _, depID := range spec.DependsOn {
depID = strings.TrimSpace(depID)
if depID == "" {
return nil, errors.New("dependency job id must not be empty")
}
if _, duplicate := seenDeps[depID]; duplicate {
continue
}
if s.state.Jobs[depID] == nil {
return nil, fmt.Errorf("dependency job %s not found", depID)
}
seenDeps[depID] = struct{}{}
cleanDeps = append(cleanDeps, depID)
}
spec.DependsOn = cleanDeps
if parentID := strings.TrimSpace(spec.ParentJobID); parentID != "" && s.state.Jobs[parentID] == nil {
return nil, fmt.Errorf("parent job %s not found", parentID)
}
now := time.Now().UTC()
maxAttempts := spec.MaxAttempts
if maxAttempts <= 0 {
maxAttempts = cfg.DefaultMaxAttempts
}
if maxAttempts <= 0 {
maxAttempts = 3
}
backoff := spec.BackoffSeconds
if backoff <= 0 {
backoff = cfg.RetryBackoffSeconds
}
if backoff <= 0 {
backoff = 15
}
maxApplyAttempts := spec.MaxApplyAttempts
if maxApplyAttempts <= 0 {
maxApplyAttempts = cfg.MasterApplyMaxAttempts
}
if maxApplyAttempts <= 0 {
maxApplyAttempts = 5
}
applyBackoff := spec.ApplyBackoffSeconds
if applyBackoff <= 0 {
applyBackoff = cfg.MasterApplyBackoffSeconds
}
if applyBackoff <= 0 {
applyBackoff = 5
}
priority := spec.Priority
if priority < -1000 {
priority = -1000
}
if priority > 1000 {
priority = 1000
}
status := "queued"
if len(spec.DependsOn) > 0 {
status = "blocked"
}
j := &core.Job{
ID: NewID("job"),
Type: spec.Type,
Payload: b,
Status: status,
Priority: priority,
ResourceClass: strings.ToLower(strings.TrimSpace(spec.ResourceClass)),
RequiredCapabilities: normalizeCapabilities(spec.RequiredCapabilities),
IdempotencyKey: strings.TrimSpace(spec.IdempotencyKey),
ParentJobID: strings.TrimSpace(spec.ParentJobID),
DependsOn: append([]string(nil), spec.DependsOn...),
MaxAttempts: maxAttempts,
BackoffSeconds: backoff,
TimeoutSeconds: spec.TimeoutSeconds,
RequiresMasterApply: spec.RequiresMasterApply,
MaxApplyAttempts: maxApplyAttempts,
ApplyBackoffSeconds: applyBackoff,
CreatedAt: now,
UpdatedAt: now,
}
s.state.Jobs[j.ID] = j
cp := *j
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) dependencyStateLocked(j *core.Job) (ready bool, terminalFailure string) {
if len(j.DependsOn) == 0 {
return true, ""
}
for _, id := range j.DependsOn {
dep := s.state.Jobs[id]
if dep == nil {
return false, ""
}
switch dep.Status {
case "done":
continue
case "failed", "canceled":
return false, "dependency " + id + " ended as " + dep.Status
default:
return false, ""
}
}
return true, ""
}
func workerMatchesJob(w WorkerHeartbeat, j *core.Job) bool {
resource := strings.ToLower(strings.TrimSpace(j.ResourceClass))
workerResource := strings.ToLower(strings.TrimSpace(w.ResourceClass))
if resource != "" && resource != "any" && resource != workerResource {
return false
}
return capabilitiesContain(w.Capabilities, j.RequiredCapabilities)
}
func (s *Store) registerWorkerLocked(h WorkerHeartbeat, now time.Time) core.WorkerState {
if s.workers == nil {
s.workers = map[string]core.WorkerState{}
}
old := s.workers[h.ID]
registered := old.RegisteredAt
if registered.IsZero() {
registered = now
}
maxc := h.MaxConcurrency
if maxc <= 0 {
maxc = 1
}
state := core.WorkerState{
ID: h.ID,
ResourceClass: strings.ToLower(strings.TrimSpace(h.ResourceClass)),
Capabilities: normalizeCapabilities(h.Capabilities),
Labels: h.Labels,
MaxConcurrency: maxc,
Version: strings.TrimSpace(h.Version),
Hostname: strings.TrimSpace(h.Hostname),
LastHeartbeat: now,
RegisteredAt: registered,
Status: "online",
}
if state.ResourceClass == "" {
state.ResourceClass = "cpu"
}
if len(state.Capabilities) == 0 {
state.Capabilities = []string{"cpu", "vector.relink"}
}
for _, j := range s.state.Jobs {
if j != nil && j.Status == "claimed" && j.ClaimedBy == h.ID && now.Before(j.LeaseUntil) {
state.Inflight++
}
}
s.workers[h.ID] = state
return state
}
func (s *Store) RegisterWorker(h WorkerHeartbeat) (core.WorkerState, error) {
h.ID = strings.TrimSpace(h.ID)
if h.ID == "" {
return core.WorkerState{}, errors.New("worker id is required")
}
s.mu.Lock()
defer s.mu.Unlock()
return s.registerWorkerLocked(h, time.Now().UTC()), nil
}
func (s *Store) HeartbeatWorker(h WorkerHeartbeat, lease time.Duration) (core.WorkerState, error) {
h.ID = strings.TrimSpace(h.ID)
if h.ID == "" {
return core.WorkerState{}, errors.New("worker id is required")
}
if lease <= 0 {
lease = 120 * time.Second
}
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UTC()
state := s.registerWorkerLocked(h, now)
for jobID, token := range h.ActiveLeases {
j := s.state.Jobs[jobID]
if j == nil || j.Status != "claimed" || j.ClaimedBy != h.ID || j.LeaseToken == "" || j.LeaseToken != token {
continue
}
// Persist a renewal only when the remaining lease is below half the
// configured lease. This gives restart-safe leases without heartbeat WAL spam.
if time.Until(j.LeaseUntil) <= lease/2 {
j.LeaseUntil = now.Add(lease)
j.UpdatedAt = now
cp := *j
if err := s.commitLocked("job.upsert", cp); err != nil {
return core.WorkerState{}, err
}
}
}
return state, nil
}
func (s *Store) normalizeJobForClaimLocked(j *core.Job, now time.Time) error {
if j == nil {
return nil
}
changed := false
if j.MaxAttempts <= 0 {
j.MaxAttempts = s.state.Config.Worker.DefaultMaxAttempts
if j.MaxAttempts <= 0 {
j.MaxAttempts = 3
}
changed = true
}
if j.Status == "claimed" && !j.LeaseUntil.IsZero() && now.After(j.LeaseUntil) {
j.ClaimedBy = ""
j.LeaseToken = ""
j.LeaseUntil = time.Time{}
if j.Attempts >= j.MaxAttempts {
j.Status = "failed"
j.Error = "lease expired after maximum attempts"
j.FinishedAt = now
} else {
j.Status = "retry_wait"
j.Error = "worker lease expired"
j.NextAttemptAt = now.Add(time.Duration(maxInt(1, j.BackoffSeconds)) * time.Second)
}
j.UpdatedAt = now
changed = true
}
if j.Status == "retry_wait" && (j.NextAttemptAt.IsZero() || !now.Before(j.NextAttemptAt)) {
j.Status = "queued"
j.NextAttemptAt = time.Time{}
j.UpdatedAt = now
changed = true
}
if j.Status == "blocked" {
ready, depFailure := s.dependencyStateLocked(j)
if depFailure != "" {
j.Status = "failed"
j.Error = depFailure
j.FinishedAt = now
j.UpdatedAt = now
changed = true
} else if ready {
j.Status = "queued"
j.UpdatedAt = now
changed = true
}
}
if changed {
cp := *j
return s.commitLocked("job.upsert", cp)
}
return nil
}
func maxInt(a, b int) int {
if a > b {
return a
}
return b
}
// ClaimJob preserves the v0.8 worker contract for compatibility. New workers
// should use ClaimJobForWorker so resource/capability routing is enforced.
func (s *Store) ClaimJob(worker string, lease time.Duration) (*core.Job, error) {
return s.ClaimJobForWorker(WorkerHeartbeat{ID: worker, ResourceClass: "cpu", Capabilities: []string{"cpu", "vector.relink"}, MaxConcurrency: 1}, lease)
}
func (s *Store) ClaimJobForWorker(h WorkerHeartbeat, lease time.Duration) (*core.Job, error) {
h.ID = strings.TrimSpace(h.ID)
if h.ID == "" {
return nil, errors.New("worker id is required")
}
if lease <= 0 {
lease = 120 * time.Second
}
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UTC()
state := s.registerWorkerLocked(h, now)
if state.Inflight >= state.MaxConcurrency {
return nil, nil
}
var chosen *core.Job
for _, j := range s.state.Jobs {
if err := s.normalizeJobForClaimLocked(j, now); err != nil {
return nil, err
}
if j == nil || j.Status != "queued" || !workerMatchesJob(h, j) {
continue
}
if chosen == nil || j.Priority > chosen.Priority || (j.Priority == chosen.Priority && j.CreatedAt.Before(chosen.CreatedAt)) {
chosen = j
}
}
if chosen == nil {
return nil, nil
}
chosen.Status = "claimed"
chosen.ClaimedBy = h.ID
chosen.LeaseToken = NewID("lease")
chosen.LeaseUntil = now.Add(lease)
chosen.Attempts++
chosen.UpdatedAt = now
if chosen.StartedAt.IsZero() {
chosen.StartedAt = now
}
cp := *chosen
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) CompleteJob(id, worker string, result json.RawMessage, jobErr string) (*core.Job, error) {
return s.CompleteJobLease(id, worker, "", result, jobErr)
}
func (s *Store) CompleteJobLease(id, worker, leaseToken string, result json.RawMessage, jobErr string) (*core.Job, error) {
s.mu.Lock()
defer s.mu.Unlock()
j, ok := s.state.Jobs[id]
if !ok || j == nil {
return nil, errors.New("job not found")
}
if j.Status != "claimed" {
return nil, fmt.Errorf("job is not claimed: %s", j.Status)
}
if j.ClaimedBy != worker {
return nil, errors.New("job claimed by another worker")
}
if j.LeaseToken != "" && leaseToken != j.LeaseToken {
return nil, errors.New("stale or invalid job lease token")
}
now := time.Now().UTC()
if !j.LeaseUntil.IsZero() && !now.Before(j.LeaseUntil) {
// A lease token fences ownership, but its validity also has a deadline.
// Reject late results even if no other worker has claimed the job yet;
// otherwise a paused partitioned worker could mutate state after expiry.
return nil, errors.New("stale or expired job lease")
}
j.Result = result
j.Error = strings.TrimSpace(jobErr)
j.UpdatedAt = now
j.ClaimedBy = ""
j.LeaseToken = ""
j.LeaseUntil = time.Time{}
if j.Error == "" {
j.NextAttemptAt = time.Time{}
if j.RequiresMasterApply {
j.Status = "apply_wait"
j.ApplyError = ""
j.ApplyNextAttemptAt = now
} else {
j.Status = "done"
j.FinishedAt = now
}
} else if j.Attempts < maxInt(1, j.MaxAttempts) {
j.Status = "retry_wait"
backoff := maxInt(1, j.BackoffSeconds)
shift := j.Attempts - 1
if shift > 6 {
shift = 6
}
delay := time.Duration(backoff*(1<<shift)) * time.Second
if delay > time.Hour {
delay = time.Hour
}
j.NextAttemptAt = now.Add(delay)
} else {
j.Status = "failed"
j.FinishedAt = now
}
cp := *j
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) Job(id string) (*core.Job, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
j := s.state.Jobs[id]
if j == nil {
return nil, false
}
cp := *j
cp.Payload = append(json.RawMessage(nil), j.Payload...)
cp.Result = append(json.RawMessage(nil), j.Result...)
return &cp, true
}
func (s *Store) JobsSnapshot(limit int, status, kind string) []core.Job {
if limit <= 0 || limit > 5000 {
limit = 500
}
status = strings.TrimSpace(status)
kind = strings.TrimSpace(kind)
s.mu.RLock()
defer s.mu.RUnlock()
out := make([]core.Job, 0)
for _, j := range s.state.Jobs {
if j == nil || (status != "" && j.Status != status) || (kind != "" && j.Type != kind) {
continue
}
cp := *j
// Admin status does not need potentially large/sensitive payloads.
cp.Payload = nil
cp.Result = nil
out = append(out, cp)
}
sort.Slice(out, func(i, j int) bool { return out[i].CreatedAt.After(out[j].CreatedAt) })
if len(out) > limit {
out = out[:limit]
}
return out
}
func (s *Store) RetryJob(id string) (*core.Job, error) {
s.mu.Lock()
defer s.mu.Unlock()
j := s.state.Jobs[id]
if j == nil {
return nil, errors.New("job not found")
}
if j.Status != "failed" && j.Status != "canceled" {
return nil, fmt.Errorf("job %s cannot be retried from status %s", id, j.Status)
}
j.Status = "queued"
j.Error = ""
j.Result = nil
j.Attempts = 0
j.NextAttemptAt = time.Time{}
j.FinishedAt = time.Time{}
j.UpdatedAt = time.Now().UTC()
cp := *j
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) WorkersSnapshot() []core.WorkerState {
s.mu.RLock()
defer s.mu.RUnlock()
now := time.Now().UTC()
staleAfter := time.Duration(s.state.Config.Worker.StaleAfterSeconds) * time.Second
if staleAfter <= 0 {
staleAfter = 60 * time.Second
}
out := make([]core.WorkerState, 0, len(s.workers))
for _, w := range s.workers {
cp := w
cp.Inflight = 0
for _, j := range s.state.Jobs {
if j != nil && j.Status == "claimed" && j.ClaimedBy == w.ID && now.Before(j.LeaseUntil) {
cp.Inflight++
}
}
if now.Sub(cp.LastHeartbeat) > staleAfter {
cp.Status = "stale"
}
out = append(out, cp)
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return out
}
func (s *Store) HasLiveWorker(resource string, capabilities ...string) bool {
resource = strings.ToLower(strings.TrimSpace(resource))
for _, w := range s.WorkersSnapshot() {
if w.Status != "online" || (resource != "" && resource != "any" && w.ResourceClass != resource) {
continue
}
if capabilitiesContain(w.Capabilities, capabilities) && w.Inflight < w.MaxConcurrency {
return true
}
}
return false
}
func (s *Store) OrchestratorStatus() OrchestratorStatus {
out := OrchestratorStatus{Workers: s.WorkersSnapshot(), JobsByStatus: map[string]int{}, JobsByType: map[string]int{}, JobsByResource: map[string]int{}}
s.mu.RLock()
defer s.mu.RUnlock()
for _, j := range s.state.Jobs {
if j == nil {
continue
}
out.JobsByStatus[j.Status]++
out.JobsByType[j.Type]++
resource := j.ResourceClass
if resource == "" {
resource = "any"
}
out.JobsByResource[resource]++
switch j.Status {
case "queued":
out.Queued++
if out.OldestQueued.IsZero() || j.CreatedAt.Before(out.OldestQueued) {
out.OldestQueued = j.CreatedAt
}
case "claimed":
out.Claimed++
case "retry_wait":
out.Retrying++
case "apply_wait":
out.Applying++
case "blocked":
out.Blocked++
case "failed":
out.Failed++
case "done":
out.Done++
}
}
return out
}
func (s *Store) PendingJobCount(kind string) int {
s.mu.RLock()
defer s.mu.RUnlock()
n := 0
for _, j := range s.state.Jobs {
if j == nil || (kind != "" && j.Type != kind) {
continue
}
if j.Status == "queued" || j.Status == "claimed" || j.Status == "retry_wait" || j.Status == "blocked" || j.Status == "apply_wait" {
n++
}
}
return n
}
// PendingMasterApplyJobs returns durable worker results whose mutation still
// needs to be committed by the master. The result payload remains persisted,
// so a master restart can resume without recomputing the worker job.
func (s *Store) PendingMasterApplyJobs(limit int) []core.Job {
if limit <= 0 || limit > 1024 {
limit = 64
}
now := time.Now().UTC()
s.mu.RLock()
defer s.mu.RUnlock()
out := make([]core.Job, 0, limit)
for _, j := range s.state.Jobs {
if j == nil || j.Status != "apply_wait" || (!j.ApplyNextAttemptAt.IsZero() && now.Before(j.ApplyNextAttemptAt)) {
continue
}
cp := *j
cp.Payload = append(json.RawMessage(nil), j.Payload...)
cp.Result = append(json.RawMessage(nil), j.Result...)
out = append(out, cp)
}
sort.Slice(out, func(i, j int) bool {
if out[i].Priority != out[j].Priority {
return out[i].Priority > out[j].Priority
}
return out[i].UpdatedAt.Before(out[j].UpdatedAt)
})
if len(out) > limit {
out = out[:limit]
}
return out
}
// FinishMasterApply closes the second phase of a worker job. Failed apply
// attempts are retried from the already persisted Result; successful worker
// computation is never lost merely because the master restarted mid-commit.
func (s *Store) FinishMasterApply(id string, applyErr error) (*core.Job, error) {
s.mu.Lock()
defer s.mu.Unlock()
j := s.state.Jobs[id]
if j == nil {
return nil, errors.New("job not found")
}
if j.Status != "apply_wait" {
return nil, fmt.Errorf("job is not awaiting master apply: %s", j.Status)
}
now := time.Now().UTC()
j.ApplyAttempts++
j.UpdatedAt = now
if applyErr == nil {
j.Status = "done"
j.ApplyError = ""
j.ApplyNextAttemptAt = time.Time{}
j.FinishedAt = now
} else {
j.ApplyError = strings.TrimSpace(applyErr.Error())
maxAttempts := j.MaxApplyAttempts
if maxAttempts <= 0 {
maxAttempts = maxInt(1, s.state.Config.Worker.MasterApplyMaxAttempts)
}
if j.ApplyAttempts >= maxAttempts {
j.Status = "failed"
j.Error = "master apply failed: " + j.ApplyError
j.FinishedAt = now
j.ApplyNextAttemptAt = time.Time{}
} else {
backoff := j.ApplyBackoffSeconds
if backoff <= 0 {
backoff = maxInt(1, s.state.Config.Worker.MasterApplyBackoffSeconds)
}
shift := j.ApplyAttempts - 1
if shift > 6 {
shift = 6
}
delay := time.Duration(backoff*(1<<shift)) * time.Second
if delay > time.Hour {
delay = time.Hour
}
j.ApplyNextAttemptAt = now.Add(delay)
}
}
cp := *j
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) CancelJob(id, reason string) error {
s.mu.Lock()
defer s.mu.Unlock()
j := s.state.Jobs[id]
if j == nil {
return errors.New("job not found")
}
if j.Status == "done" || j.Status == "failed" || j.Status == "canceled" {
return nil
}
now := time.Now().UTC()
j.Status = "canceled"
j.Error = strings.TrimSpace(reason)
j.ClaimedBy = ""
j.LeaseToken = ""
j.LeaseUntil = time.Time{}
j.FinishedAt = now
j.UpdatedAt = now
cp := *j
return s.commitLocked("job.upsert", cp)
}
// PruneTerminalJobs bounds durable scheduler history while preserving jobs that
// are still referenced as dependencies by non-terminal work.
func (s *Store) PruneTerminalJobs(retention time.Duration, maxTerminal int) (int, error) {
if retention <= 0 {
retention = 7 * 24 * time.Hour
}
if maxTerminal < 100 {
maxTerminal = 20000
}
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UTC()
cutoff := now.Add(-retention)
referenced := map[string]bool{}
for _, j := range s.state.Jobs {
if j == nil || j.Status == "done" || j.Status == "failed" || j.Status == "canceled" {
continue
}
for _, dep := range j.DependsOn {
referenced[dep] = true
}
}
type terminal struct {
id string
t time.Time
}
items := make([]terminal, 0)
for id, j := range s.state.Jobs {
if j == nil || referenced[id] || (j.Status != "done" && j.Status != "failed" && j.Status != "canceled") {
continue
}
t := j.FinishedAt
if t.IsZero() {
t = j.UpdatedAt
}
items = append(items, terminal{id: id, t: t})
}
sort.Slice(items, func(i, j int) bool { return items[i].t.Before(items[j].t) })
remove := map[string]bool{}
for _, x := range items {
if !x.t.IsZero() && x.t.Before(cutoff) {
remove[x.id] = true
}
}
remaining := len(items) - len(remove)
if remaining > maxTerminal {
need := remaining - maxTerminal
for _, x := range items {
if need == 0 {
break
}
if remove[x.id] {
continue
}
remove[x.id] = true
need--
}
}
if len(remove) == 0 {
return 0, nil
}
ids := make([]string, 0, len(remove))
for id := range remove {
delete(s.state.Jobs, id)
ids = append(ids, id)
}
sort.Strings(ids)
return len(ids), s.commitLocked("job.delete", ids)
}
// IsOrchestratorLeader prevents two clustered masters from planning/applying
// the same durable work. Standalone installations are always authoritative.
func (s *Store) IsOrchestratorLeader() bool {
s.mu.RLock()
defer s.mu.RUnlock()
cfg := s.state.Config.Cluster
if !cfg.Enabled {
return true
}
leader := s.state.Cluster.LeaderID
if !cfg.AutoElection && leader == "" {
leader = cfg.LeaderID
}
return leader != "" && leader == cfg.NodeID && s.state.Cluster.Role == ClusterLeader
}
@@ -0,0 +1,294 @@
package store
import (
"encoding/json"
"strings"
"testing"
"time"
)
func TestOrchestratorCapabilityRoutingAndLeaseFencing(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
j, err := s.EnqueueJobSpec(JobSpec{
Type: "model.chat", Payload: map[string]any{"input": "x"}, Priority: 100,
ResourceClass: "gpu", RequiredCapabilities: []string{"gpu", "model.chat"},
})
if err != nil {
t.Fatal(err)
}
cpu := WorkerHeartbeat{ID: "cpu-1", ResourceClass: "cpu", Capabilities: []string{"cpu", "vector.relink"}, MaxConcurrency: 1}
got, err := s.ClaimJobForWorker(cpu, time.Minute)
if err != nil {
t.Fatal(err)
}
if got != nil {
t.Fatalf("CPU worker must not claim GPU job: %+v", got)
}
gpu := WorkerHeartbeat{ID: "gpu-1", ResourceClass: "gpu", Capabilities: []string{"gpu", "model.chat"}, MaxConcurrency: 1}
got, err = s.ClaimJobForWorker(gpu, time.Minute)
if err != nil {
t.Fatal(err)
}
if got == nil || got.ID != j.ID || got.LeaseToken == "" {
t.Fatalf("GPU claim=%+v", got)
}
if _, err := s.CompleteJobLease(got.ID, gpu.ID, "stale-token", json.RawMessage(`{"ok":true}`), ""); err == nil || !strings.Contains(err.Error(), "lease") {
t.Fatalf("stale lease must be fenced, err=%v", err)
}
done, err := s.CompleteJobLease(got.ID, gpu.ID, got.LeaseToken, json.RawMessage(`{"ok":true}`), "")
if err != nil {
t.Fatal(err)
}
if done.Status != "done" {
t.Fatalf("status=%s", done.Status)
}
}
func TestOrchestratorRetryDependencyAndWorkerStaleness(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
cfg := s.Config()
cfg.Worker.RetryBackoffSeconds = 1
cfg.Worker.HeartbeatSeconds = 2
cfg.Worker.StaleAfterSeconds = 2
if err := s.UpdateConfig(cfg); err != nil {
t.Fatal(err)
}
parent, err := s.EnqueueJobSpec(JobSpec{Type: "parent", Payload: map[string]any{}, ResourceClass: "cpu", MaxAttempts: 2, BackoffSeconds: 1})
if err != nil {
t.Fatal(err)
}
child, err := s.EnqueueJobSpec(JobSpec{Type: "child", Payload: map[string]any{}, ResourceClass: "cpu", DependsOn: []string{parent.ID}})
if err != nil {
t.Fatal(err)
}
if child.Status != "blocked" {
t.Fatalf("child status=%s", child.Status)
}
w := WorkerHeartbeat{ID: "cpu-1", ResourceClass: "cpu", Capabilities: []string{"cpu"}, MaxConcurrency: 1}
p1, err := s.ClaimJobForWorker(w, time.Minute)
if err != nil {
t.Fatal(err)
}
if p1 == nil || p1.ID != parent.ID {
t.Fatalf("claim=%+v", p1)
}
retry, err := s.CompleteJobLease(p1.ID, w.ID, p1.LeaseToken, nil, "temporary")
if err != nil {
t.Fatal(err)
}
if retry.Status != "retry_wait" || retry.NextAttemptAt.IsZero() {
t.Fatalf("retry=%+v", retry)
}
// Avoid sleeping: advance this job's retry clock under the package lock.
s.mu.Lock()
s.state.Jobs[parent.ID].NextAttemptAt = time.Now().Add(-time.Second)
s.mu.Unlock()
p2, err := s.ClaimJobForWorker(w, time.Minute)
if err != nil {
t.Fatal(err)
}
if p2 == nil || p2.ID != parent.ID || p2.Attempts != 2 {
t.Fatalf("second claim=%+v", p2)
}
if _, err := s.CompleteJobLease(p2.ID, w.ID, p2.LeaseToken, json.RawMessage(`{"ok":true}`), ""); err != nil {
t.Fatal(err)
}
c, err := s.ClaimJobForWorker(w, time.Minute)
if err != nil {
t.Fatal(err)
}
if c == nil || c.ID != child.ID {
t.Fatalf("child claim after dependency=%+v", c)
}
s.mu.Lock()
ws := s.workers[w.ID]
ws.LastHeartbeat = time.Now().Add(-2 * time.Second)
s.workers[w.ID] = ws
s.mu.Unlock()
workers := s.WorkersSnapshot()
if len(workers) != 1 || workers[0].Status != "stale" {
t.Fatalf("workers=%+v", workers)
}
}
func TestDependencyFailureFailsBlockedChild(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
p, err := s.EnqueueJobSpec(JobSpec{Type: "parent", Payload: map[string]any{}, ResourceClass: "cpu", MaxAttempts: 1})
if err != nil {
t.Fatal(err)
}
c, err := s.EnqueueJobSpec(JobSpec{Type: "child", Payload: map[string]any{}, ResourceClass: "cpu", DependsOn: []string{p.ID}})
if err != nil {
t.Fatal(err)
}
w := WorkerHeartbeat{ID: "cpu", ResourceClass: "cpu", Capabilities: []string{"cpu"}, MaxConcurrency: 1}
pj, _ := s.ClaimJobForWorker(w, time.Minute)
if _, err := s.CompleteJobLease(pj.ID, w.ID, pj.LeaseToken, nil, "permanent"); err != nil {
t.Fatal(err)
}
got, err := s.ClaimJobForWorker(w, time.Minute)
if err != nil {
t.Fatal(err)
}
if got != nil {
t.Fatalf("failed dependency child must not be claimable: %+v", got)
}
cj, ok := s.Job(c.ID)
if !ok || cj.Status != "failed" || !strings.Contains(cj.Error, "dependency") {
t.Fatalf("child=%+v", cj)
}
}
func TestPruneTerminalJobsPreservesReferencedDependencies(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
p, err := s.EnqueueJobSpec(JobSpec{Type: "parent", Payload: map[string]any{}, ResourceClass: "cpu"})
if err != nil {
t.Fatal(err)
}
w := WorkerHeartbeat{ID: "cpu", ResourceClass: "cpu", Capabilities: []string{"cpu"}, MaxConcurrency: 1}
pj, _ := s.ClaimJobForWorker(w, time.Minute)
if _, err := s.CompleteJobLease(pj.ID, w.ID, pj.LeaseToken, nil, ""); err != nil {
t.Fatal(err)
}
_, err = s.EnqueueJobSpec(JobSpec{Type: "child", Payload: map[string]any{}, ResourceClass: "gpu", RequiredCapabilities: []string{"gpu"}, DependsOn: []string{p.ID}})
if err != nil {
t.Fatal(err)
}
s.mu.Lock()
s.state.Jobs[p.ID].FinishedAt = time.Now().Add(-48 * time.Hour)
s.state.Jobs[p.ID].UpdatedAt = s.state.Jobs[p.ID].FinishedAt
s.mu.Unlock()
n, err := s.PruneTerminalJobs(time.Hour, 100)
if err != nil {
t.Fatal(err)
}
if n != 0 {
t.Fatalf("pruned referenced dependency: %d", n)
}
if _, ok := s.Job(p.ID); !ok {
t.Fatal("referenced parent was removed")
}
}
func TestExpiredLeaseCannotCompleteBeforeReclaim(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
_, err = s.EnqueueJobSpec(JobSpec{Type: "cpu.task", Payload: map[string]any{"x": 1}, ResourceClass: "cpu", RequiredCapabilities: []string{"cpu"}})
if err != nil {
t.Fatal(err)
}
w := WorkerHeartbeat{ID: "cpu-1", ResourceClass: "cpu", Capabilities: []string{"cpu"}, MaxConcurrency: 1}
j, err := s.ClaimJobForWorker(w, time.Minute)
if err != nil || j == nil {
t.Fatalf("claim=%+v err=%v", j, err)
}
s.mu.Lock()
s.state.Jobs[j.ID].LeaseUntil = time.Now().Add(-time.Second)
s.mu.Unlock()
if _, err := s.CompleteJobLease(j.ID, w.ID, j.LeaseToken, json.RawMessage(`{"ok":true}`), ""); err == nil || !strings.Contains(err.Error(), "expired") {
t.Fatalf("expired lease completion must be fenced, err=%v", err)
}
cur, ok := s.Job(j.ID)
if !ok || cur.Status != "claimed" {
t.Fatalf("late completion mutated job before scheduler recovery: %+v", cur)
}
// Any subsequent scheduler/claim pass normalizes the expired lease and may
// make the job retryable according to its durable attempt policy.
other := WorkerHeartbeat{ID: "cpu-2", ResourceClass: "cpu", Capabilities: []string{"cpu"}, MaxConcurrency: 1}
next, err := s.ClaimJobForWorker(other, time.Minute)
if err != nil {
t.Fatal(err)
}
if next != nil {
t.Fatalf("retry backoff should prevent immediate reclaim: %+v", next)
}
cur, _ = s.Job(j.ID)
if cur.Status != "retry_wait" || cur.Error != "worker lease expired" {
t.Fatalf("expired lease not normalized to retry_wait: %+v", cur)
}
}
func TestOrchestratorRejectsDanglingDependencies(t *testing.T) {
s, err := New(t.TempDir())
if err != nil {
t.Fatal(err)
}
defer s.Close()
if _, err := s.EnqueueJobSpec(JobSpec{Type: "child", Payload: map[string]any{}, DependsOn: []string{"job_missing"}}); err == nil || !strings.Contains(err.Error(), "not found") {
t.Fatalf("dangling dependency must fail closed, err=%v", err)
}
if s.PendingJobCount("") != 0 {
t.Fatalf("failed planning left durable work behind: %+v", s.OrchestratorStatus())
}
}
func TestApplyWaitJobSurvivesRestart(t *testing.T) {
dir := t.TempDir()
s, err := New(dir)
if err != nil {
t.Fatal(err)
}
j, err := s.EnqueueJobSpec(JobSpec{Type: "master.apply", Payload: map[string]any{"x": 1}, ResourceClass: "cpu", RequiredCapabilities: []string{"cpu"}, RequiresMasterApply: true})
if err != nil {
t.Fatal(err)
}
w := WorkerHeartbeat{ID: "cpu-1", ResourceClass: "cpu", Capabilities: []string{"cpu"}, MaxConcurrency: 1}
claimed, err := s.ClaimJobForWorker(w, time.Minute)
if err != nil || claimed == nil || claimed.ID != j.ID {
t.Fatalf("claim=%+v err=%v", claimed, err)
}
result := json.RawMessage(`{"ok":true}`)
completed, err := s.CompleteJobLease(claimed.ID, w.ID, claimed.LeaseToken, result, "")
if err != nil || completed.Status != "apply_wait" {
t.Fatalf("complete=%+v err=%v", completed, err)
}
if err := s.Close(); err != nil {
t.Fatal(err)
}
s2, err := New(dir)
if err != nil {
t.Fatal(err)
}
defer s2.Close()
recovered, ok := s2.Job(j.ID)
if !ok || recovered.Status != "apply_wait" {
t.Fatalf("durable apply_wait job not recovered: %+v", recovered)
}
if len(recovered.Payload) == 0 || string(recovered.Result) != string(result) {
t.Fatalf("recovered payload/result incomplete: %+v", recovered)
}
pending := s2.PendingMasterApplyJobs(10)
if len(pending) != 1 || pending[0].ID != j.ID {
t.Fatalf("recovered master apply queue=%+v", pending)
}
}
+94 -132
View File
@@ -45,6 +45,8 @@ type Store struct {
clusterLogMu sync.Mutex
clusterLog *ClusterLog
provenanceSourceIDs map[string]map[string]struct{}
workers map[string]core.WorkerState
synapseAdj map[string]map[string]*core.Synapse
}
func New(dir string) (*Store, error) {
@@ -54,7 +56,7 @@ func New(dir string) (*Store, error) {
if err := os.MkdirAll(dir, 0700); err != nil {
return nil, err
}
s := &Store{dir: dir, indexes: map[int]*vector.HNSW{}, diskIndexes: map[int]*vector.PQIndex{}, provenanceSourceIDs: map[string]map[string]struct{}{}}
s := &Store{dir: dir, indexes: map[int]*vector.HNSW{}, diskIndexes: map[int]*vector.PQIndex{}, provenanceSourceIDs: map[string]map[string]struct{}{}, workers: map[string]core.WorkerState{}, synapseAdj: map[string]map[string]*core.Synapse{}}
s.state = core.PersistedState{Config: core.DefaultConfig(), Memories: map[string]*core.Memory{}, Synapses: map[string]*core.Synapse{}, Jobs: map[string]*core.Job{}, Goals: map[string]*core.Goal{}, Sources: map[string]*core.KnowledgeSource{}, ResearchRuns: map[string]*core.ResearchRun{}}
_ = s.loadJSON(filepath.Join(dir, "state.json"), &s.state)
_ = s.loadJSON(filepath.Join(dir, "secrets.json"), &s.secrets)
@@ -111,6 +113,7 @@ func New(dir string) (*Store, error) {
if err := s.replayWAL(); err != nil {
return nil, err
}
s.rebuildSynapseAdjLocked()
applyNewDefaults(&s.state.Config)
if s.state.Cluster.Term < s.state.Config.Cluster.Term {
s.state.Cluster.Term = s.state.Config.Cluster.Term
@@ -235,6 +238,18 @@ func applyNewDefaults(c *core.Config) {
if c.Brain.TypeWeights == nil {
c.Brain.TypeWeights = d.Brain.TypeWeights
}
if c.Brain.GraphMaxHops == 0 {
c.Brain.GraphMaxHops = d.Brain.GraphMaxHops
}
if c.Brain.GraphHopDecay == 0 {
c.Brain.GraphHopDecay = d.Brain.GraphHopDecay
}
if c.Brain.GraphMaxExpansion == 0 {
c.Brain.GraphMaxExpansion = d.Brain.GraphMaxExpansion
}
if c.Brain.GraphMinEdgeWeight == 0 {
c.Brain.GraphMinEdgeWeight = d.Brain.GraphMinEdgeWeight
}
if c.Brain.LearningPolicy.MaxMemoryTextChars == 0 && c.Brain.LearningPolicy.DuplicateSimilarity == 0 && c.Brain.LearningPolicy.SemanticMinConfirmations == 0 {
c.Brain.LearningPolicy = d.Brain.LearningPolicy
} else {
@@ -484,6 +499,65 @@ func applyNewDefaults(c *core.Config) {
if c.Cluster.LogSegmentBytes == 0 {
c.Cluster.LogSegmentBytes = d.Cluster.LogSegmentBytes
}
if c.Worker.LeaseSeconds == 0 {
c.Worker.LeaseSeconds = d.Worker.LeaseSeconds
}
if c.Worker.HeartbeatSeconds == 0 && c.Worker.StaleAfterSeconds == 0 && c.Worker.DefaultMaxAttempts == 0 {
legacyLease := c.Worker.LeaseSeconds
c.Worker = d.Worker
if legacyLease > 0 {
c.Worker.LeaseSeconds = legacyLease
}
} else {
if c.Worker.HeartbeatSeconds == 0 {
c.Worker.HeartbeatSeconds = d.Worker.HeartbeatSeconds
}
if c.Worker.StaleAfterSeconds == 0 {
c.Worker.StaleAfterSeconds = d.Worker.StaleAfterSeconds
}
if c.Worker.DefaultMaxAttempts == 0 {
c.Worker.DefaultMaxAttempts = d.Worker.DefaultMaxAttempts
}
if c.Worker.RetryBackoffSeconds == 0 {
c.Worker.RetryBackoffSeconds = d.Worker.RetryBackoffSeconds
}
if c.Worker.MaxQueuedJobs == 0 {
c.Worker.MaxQueuedJobs = d.Worker.MaxQueuedJobs
}
if c.Worker.MasterApplyMaxAttempts == 0 {
c.Worker.MasterApplyMaxAttempts = d.Worker.MasterApplyMaxAttempts
}
if c.Worker.MasterApplyBackoffSeconds == 0 {
c.Worker.MasterApplyBackoffSeconds = d.Worker.MasterApplyBackoffSeconds
}
if c.Worker.JobRetentionHours == 0 {
c.Worker.JobRetentionHours = d.Worker.JobRetentionHours
}
if c.Worker.MaxTerminalJobs == 0 {
c.Worker.MaxTerminalJobs = d.Worker.MaxTerminalJobs
}
if c.Worker.GraphBackfillIntervalS == 0 {
c.Worker.GraphBackfillIntervalS = d.Worker.GraphBackfillIntervalS
}
if c.Worker.GraphBackfillBatchSize == 0 {
c.Worker.GraphBackfillBatchSize = d.Worker.GraphBackfillBatchSize
}
if c.Worker.GraphBackfillMaxQueued == 0 {
c.Worker.GraphBackfillMaxQueued = d.Worker.GraphBackfillMaxQueued
}
if c.Worker.GraphBackfillMinDegree == 0 {
c.Worker.GraphBackfillMinDegree = d.Worker.GraphBackfillMinDegree
}
if c.Worker.GraphCandidateMultiplier == 0 {
c.Worker.GraphCandidateMultiplier = d.Worker.GraphCandidateMultiplier
}
if c.Worker.GraphRetryAfterMinutes == 0 {
c.Worker.GraphRetryAfterMinutes = d.Worker.GraphRetryAfterMinutes
}
if c.Worker.DistributedInferenceWaitS == 0 {
c.Worker.DistributedInferenceWaitS = d.Worker.DistributedInferenceWaitS
}
}
}
func migrateMemories(memories map[string]*core.Memory, localShard string) {
@@ -1010,51 +1084,7 @@ func (s *Store) searchVectorLocked(q []float32, k int, min float64, graphBonus f
if len(hits) > k {
hits = hits[:k]
}
if graphBonus > 0 && len(hits) > 0 {
base := map[string]float64{}
for _, h := range hits {
base[h.Memory.ID] = h.Score
}
for _, syn := range s.state.Synapses {
var to string
if _, ok := base[syn.A]; ok {
to = syn.B
} else if _, ok := base[syn.B]; ok {
to = syn.A
} else {
continue
}
meta, ok := s.state.Memories[to]
if !ok || !memorySearchable(meta) {
continue
}
m, fullOK := s.fullMemoryForReadLocked(to)
if !fullOK || len(m.Vector) != len(q) {
continue
}
bonus := graphBonus * syn.Weight
found := false
for i := range hits {
if hits[i].Memory.ID == to {
hits[i].GraphBoost += bonus
hits[i].Score += bonus
found = true
break
}
}
if !found && len(hits) < k {
sim := vector.Cosine(q, m.Vector)
if sim >= min {
baseScore := sim * 0.5
hits = append(hits, SearchHit{Memory: cloneMemory(m), Similarity: sim, BaseScore: baseScore, GraphBoost: bonus, Score: bonus + baseScore, TypeWeight: 1, SalienceFactor: 1, ConfidenceFactor: 1, CandidateSource: "synapse"})
}
}
}
sort.Slice(hits, func(i, j int) bool { return hits[i].Score > hits[j].Score })
if len(hits) > k {
hits = hits[:k]
}
}
hits = s.graphExpandLocked(q, hits, k, min, graphBonus)
return hits
}
@@ -1066,32 +1096,7 @@ func edgeKey(a, b string) string {
}
func (s *Store) Reinforce(a, b string, similarity, delta, decayPerDay, maxWeight float64) error {
if a == "" || b == "" || a == b {
return nil
}
s.mu.Lock()
defer s.mu.Unlock()
if s.state.Memories[a] == nil || s.state.Memories[b] == nil {
return nil
}
key := edgeKey(a, b)
now := time.Now().UTC()
syn, ok := s.state.Synapses[key]
if !ok {
syn = &core.Synapse{A: a, B: b, Similarity: similarity, LastUpdated: now}
s.state.Synapses[key] = syn
}
days := now.Sub(syn.LastUpdated).Hours() / 24
if days > 0 && decayPerDay > 0 {
syn.Weight *= pow(1-decayPerDay, days)
}
syn.Weight = vector.Clamp(syn.Weight+delta, -maxWeight, maxWeight)
if similarity > syn.Similarity {
syn.Similarity = similarity
}
syn.Activations++
syn.LastUpdated = now
return s.commitLocked("synapse.upsert", *syn)
return s.ReinforceRelation(a, b, "association", similarity, delta, decayPerDay, maxWeight)
}
func pow(base, exp float64) float64 {
@@ -1113,6 +1118,7 @@ func (s *Store) DecayAndPruneSynapses(decayPerDay, pruneBelow float64) (int, err
syn.LastUpdated = now
}
if pruneBelow > 0 && math.Abs(syn.Weight) < pruneBelow {
s.unindexSynapseLocked(syn)
delete(s.state.Synapses, key)
pruned++
}
@@ -1348,66 +1354,6 @@ func (s *Store) RecentUsage(limit int) []core.UsageEvent {
return out
}
func (s *Store) EnqueueJob(kind string, payload any) (*core.Job, error) {
b, err := json.Marshal(payload)
if err != nil {
return nil, err
}
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UTC()
j := &core.Job{ID: NewID("job"), Type: kind, Payload: b, Status: "queued", CreatedAt: now, UpdatedAt: now}
s.state.Jobs[j.ID] = j
return j, s.commitLocked("job.upsert", *j)
}
func (s *Store) ClaimJob(worker string, lease time.Duration) (*core.Job, error) {
s.mu.Lock()
defer s.mu.Unlock()
now := time.Now().UTC()
var chosen *core.Job
for _, j := range s.state.Jobs {
if j.Status == "claimed" && !j.LeaseUntil.IsZero() && now.After(j.LeaseUntil) {
j.Status = "queued"
j.ClaimedBy = ""
}
if j.Status == "queued" && (chosen == nil || j.CreatedAt.Before(chosen.CreatedAt)) {
chosen = j
}
}
if chosen == nil {
return nil, nil
}
chosen.Status = "claimed"
chosen.ClaimedBy = worker
chosen.LeaseUntil = now.Add(lease)
chosen.UpdatedAt = now
cp := *chosen
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) CompleteJob(id, worker string, result json.RawMessage, jobErr string) (*core.Job, error) {
s.mu.Lock()
defer s.mu.Unlock()
j, ok := s.state.Jobs[id]
if !ok {
return nil, errors.New("job not found")
}
if j.ClaimedBy != worker {
return nil, errors.New("job claimed by another worker")
}
j.Result = result
j.Error = jobErr
j.UpdatedAt = time.Now().UTC()
if jobErr != "" {
j.Status = "failed"
} else {
j.Status = "done"
}
cp := *j
return &cp, s.commitLocked("job.upsert", cp)
}
func (s *Store) DeleteMemory(id string) error {
s.mu.Lock()
defer s.mu.Unlock()
@@ -1421,6 +1367,7 @@ func (s *Store) DeleteMemory(id string) error {
}
for k, x := range s.state.Synapses {
if x.A == id || x.B == id {
s.unindexSynapseLocked(x)
delete(s.state.Synapses, k)
}
}
@@ -1459,6 +1406,9 @@ func (s *Store) validateConfigLocked(c core.Config) error {
if c.Brain.MinSimilarity < -1 || c.Brain.MinSimilarity > 1 {
return errors.New("brain.min_similarity must be -1..1")
}
if c.Brain.GraphMaxHops < 1 || c.Brain.GraphMaxHops > 6 || c.Brain.GraphHopDecay <= 0 || c.Brain.GraphHopDecay > 1 || c.Brain.GraphMaxExpansion < 1 || c.Brain.GraphMaxExpansion > 5000 || c.Brain.GraphMinEdgeWeight < 0 || c.Brain.GraphMinEdgeWeight > c.Brain.MaxSynapseWeight {
return errors.New("invalid brain graph traversal configuration")
}
if c.Brain.Index.Enabled {
mode := indexMode(c)
if c.Brain.Index.Mode != "" && c.Brain.Index.Mode != "hnsw" && c.Brain.Index.Mode != "hybrid" && c.Brain.Index.Mode != "disk-pq" {
@@ -1656,6 +1606,18 @@ func (s *Store) validateConfigLocked(c core.Config) error {
return fmt.Errorf("cluster.term %d cannot be lower than persisted term %d", c.Cluster.Term, s.state.Cluster.Term)
}
}
if c.Worker.LeaseSeconds < 10 || c.Worker.LeaseSeconds > 3600 || c.Worker.HeartbeatSeconds < 2 || c.Worker.HeartbeatSeconds >= c.Worker.LeaseSeconds || c.Worker.StaleAfterSeconds < c.Worker.HeartbeatSeconds || c.Worker.StaleAfterSeconds > 7200 {
return errors.New("invalid worker lease/heartbeat/stale timing")
}
if c.Worker.DefaultMaxAttempts < 1 || c.Worker.DefaultMaxAttempts > 20 || c.Worker.RetryBackoffSeconds < 1 || c.Worker.RetryBackoffSeconds > 3600 || c.Worker.MaxQueuedJobs < 16 || c.Worker.MaxQueuedJobs > 1000000 || c.Worker.MasterApplyMaxAttempts < 1 || c.Worker.MasterApplyMaxAttempts > 20 || c.Worker.MasterApplyBackoffSeconds < 1 || c.Worker.MasterApplyBackoffSeconds > 3600 || c.Worker.JobRetentionHours < 1 || c.Worker.JobRetentionHours > 8760 || c.Worker.MaxTerminalJobs < 100 || c.Worker.MaxTerminalJobs > 1000000 {
return errors.New("invalid worker retry/queue/retention configuration")
}
if c.Worker.GraphBackfillIntervalS < 2 || c.Worker.GraphBackfillIntervalS > 3600 || c.Worker.GraphBackfillBatchSize < 1 || c.Worker.GraphBackfillBatchSize > 4096 || c.Worker.GraphBackfillMaxQueued < 1 || c.Worker.GraphBackfillMaxQueued > c.Worker.MaxQueuedJobs || c.Worker.GraphBackfillMinDegree < 1 || c.Worker.GraphBackfillMinDegree > 100 || c.Worker.GraphCandidateMultiplier < 2 || c.Worker.GraphCandidateMultiplier > 64 || c.Worker.GraphRetryAfterMinutes < 1 || c.Worker.GraphRetryAfterMinutes > 43200 {
return errors.New("invalid worker graph backfill configuration")
}
if c.Worker.DistributedInferenceWaitS < 5 || c.Worker.DistributedInferenceWaitS > 3600 {
return errors.New("worker.distributed_inference_wait_seconds must be 5..3600")
}
if c.OpenAI.DailyBudgetUSD < 0 || c.OpenAI.MonthlyBudgetUSD < 0 {
return errors.New("budgets must be >= 0")
}
+10
View File
@@ -209,6 +209,7 @@ func (s *Store) applyWALEvent(ev walEvent) error {
return err
}
s.state.Synapses[edgeKey(syn.A, syn.B)] = &syn
s.indexSynapseLocked(&syn)
case "synapse.replace":
var items []core.Synapse
if err := json.Unmarshal(ev.Data, &items); err != nil {
@@ -219,6 +220,7 @@ func (s *Store) applyWALEvent(ev walEvent) error {
x := items[i]
s.state.Synapses[edgeKey(x.A, x.B)] = &x
}
s.rebuildSynapseAdjLocked()
case "usage.add":
var x core.UsageEvent
if err := json.Unmarshal(ev.Data, &x); err != nil {
@@ -231,6 +233,14 @@ func (s *Store) applyWALEvent(ev walEvent) error {
return err
}
s.state.Jobs[x.ID] = &x
case "job.delete":
var ids []string
if err := json.Unmarshal(ev.Data, &ids); err != nil {
return err
}
for _, id := range ids {
delete(s.state.Jobs, id)
}
case "maintenance.set":
return json.Unmarshal(ev.Data, &s.state.Maintenance)
case "goal.upsert":
+89 -2
View File
@@ -16,6 +16,10 @@ components:
WorkerKey:
type: http
scheme: bearer
ControlReadKey:
type: http
scheme: bearer
description: Scoped read token used by the Control Center for graph/orchestrator status.
AdminToken:
type: apiKey
in: header
@@ -1324,13 +1328,46 @@ paths:
responses:
'200':
description: Open truth-key conflicts
/api/v1/integrations/graph/status:
get:
security:
- ControlReadKey: []
- AppKey: []
- AdminToken: []
description: Structural knowledge-graph health and graph/backfill configuration.
responses:
'200':
description: Graph statistics including linked/isolated/multi-linked nodes and connected components.
/api/v1/integrations/orchestrator/status:
get:
security:
- ControlReadKey: []
- AppKey: []
- AdminToken: []
description: Read-only master/subagent scheduler status, workers, queues and offload settings.
responses:
'200': {description: Orchestrator status}
/api/v1/worker/register:
post:
security:
- WorkerKey: []
description: Register or refresh a capability-bearing CPU/GPU subagent.
responses:
'200': {description: Worker state}
/api/v1/worker/heartbeat:
post:
security:
- WorkerKey: []
description: Refresh worker liveness and renew matching active leases using lease-token fencing.
responses:
'200': {description: Worker state}
/api/v1/worker/claim:
post:
security:
- WorkerKey: []
responses:
'200':
description: Claimed CPU job
description: Highest-priority compatible leased job for this worker resource/capability set.
'204':
description: No job available
/api/v1/worker/complete:
@@ -1339,7 +1376,9 @@ paths:
- WorkerKey: []
responses:
'200':
description: Job completion accepted
description: Lease-fenced result accepted. Master-state mutations may enter durable apply_wait before done.
'409':
description: Stale/expired lease, wrong worker, or job no longer claimable.
/internal/v1/cluster/request-vote:
post:
security:
@@ -1554,6 +1593,54 @@ paths:
responses:
'200':
description: Secrets updated
/admin/api/graph/status:
get:
security: [{AdminToken: []}]
responses:
'200': {description: Graph and orchestrator status with effective graph configuration}
/admin/api/graph/backfill:
post:
security: [{AdminToken: []}]
responses:
'202': {description: Bounded graph-backfill planning cycle accepted}
/admin/api/orchestrator/status:
get:
security: [{AdminToken: []}]
responses:
'200': {description: Master/subagent scheduler status}
/admin/api/orchestrator/jobs:
get:
security: [{AdminToken: []}]
parameters:
- {name: limit, in: query, schema: {type: integer, minimum: 1, maximum: 5000}}
- {name: status, in: query, schema: {type: string}}
- {name: type, in: query, schema: {type: string}}
responses:
'200': {description: Bounded job list without large payload/result bodies}
/admin/api/orchestrator/jobs/{id}/retry:
post:
security: [{AdminToken: []}]
parameters:
- {name: id, in: path, required: true, schema: {type: string}}
responses:
'200': {description: Failed/canceled job reset to queued}
'409': {description: Job is not retryable from its current state}
/admin/api/orchestrator/jobs/{id}/cancel:
post:
security: [{AdminToken: []}]
parameters:
- {name: id, in: path, required: true, schema: {type: string}}
requestBody:
required: false
content:
application/json:
schema:
type: object
properties:
reason: {type: string}
responses:
'200': {description: Non-terminal job canceled and any active lease fenced}
'404': {description: Job not found}
/admin/api/provider-health:
post:
security: